diff --git a/lib/hypervisor/firecracker/binaries.go b/lib/hypervisor/firecracker/binaries.go index 52fedd90a..32389d13a 100644 --- a/lib/hypervisor/firecracker/binaries.go +++ b/lib/hypervisor/firecracker/binaries.go @@ -84,19 +84,23 @@ func resolveBinaryPath(p *paths.Paths, version string) (string, error) { return "", fmt.Errorf("paths are required when using embedded firecracker binaries") } - return extractBinary(p, parseVersion(version)) + parsedVersion, err := parseVersion(version) + if err != nil { + return "", err + } + return extractBinary(p, parsedVersion) } -func parseVersion(version string) Version { +func parseVersion(version string) (Version, error) { if version == "" { - return defaultVersion + return defaultVersion, nil } for _, supported := range supportedVersions { if version == string(supported) { - return supported + return supported, nil } } - return defaultVersion + return "", fmt.Errorf("unsupported firecracker version %q", version) } func extractBinary(p *paths.Paths, version Version) (string, error) { diff --git a/lib/hypervisor/firecracker/binaries_test.go b/lib/hypervisor/firecracker/binaries_test.go index 6683fed31..2ebfe936e 100644 --- a/lib/hypervisor/firecracker/binaries_test.go +++ b/lib/hypervisor/firecracker/binaries_test.go @@ -35,12 +35,21 @@ func TestResolveBinaryPathInvalidCustomPath(t *testing.T) { assert.Contains(t, err.Error(), "invalid firecracker custom binary path") } -func TestParseVersionFallback(t *testing.T) { - assert.Equal(t, V1_16_1, defaultVersion) - assert.Equal(t, defaultVersion, parseVersion("")) - assert.Equal(t, defaultVersion, parseVersion("unknown")) - assert.Equal(t, V1_14_2, parseVersion("v1.14.2")) - assert.Equal(t, V1_16_1, parseVersion("v1.16.1")) +func TestParseVersion(t *testing.T) { + version, err := parseVersion("") + require.NoError(t, err) + assert.Equal(t, defaultVersion, version) + + _, err = parseVersion("unknown") + require.ErrorContains(t, err, "unsupported firecracker version") + + version, err = parseVersion("v1.14.2") + require.NoError(t, err) + assert.Equal(t, V1_14_2, version) + + version, err = parseVersion("v1.16.1") + require.NoError(t, err) + assert.Equal(t, V1_16_1, version) } func TestResolveEmbeddedBinaryVersions(t *testing.T) { diff --git a/lib/images/credentials_test.go b/lib/images/credentials_test.go index 46212ba77..7f680b980 100644 --- a/lib/images/credentials_test.go +++ b/lib/images/credentials_test.go @@ -5,7 +5,6 @@ import ( "net/http" "net/http/httptest" "os" - "path/filepath" "strings" "testing" "time" @@ -203,7 +202,7 @@ func TestRecoverInterruptedCredentialedPullFailsForFreshRetry(t *testing.T) { assert.Equal(t, ErrBorrowedCredentialsExpired.Error(), *stored.Error) assert.Zero(t, m.queue.QueueLength()) - data, err := os.ReadFile(filepath.Join(p.ImageDigestDir(repository, strings.TrimPrefix(digest, "sha256:")), "metadata.json")) + data, err := os.ReadFile(p.ImageContentMetadata(strings.TrimPrefix(digest, "sha256:"))) require.NoError(t, err) assert.NotContains(t, string(data), "password") } diff --git a/lib/images/disk_usage.go b/lib/images/disk_usage.go index 47914c27b..b65dba1ac 100644 --- a/lib/images/disk_usage.go +++ b/lib/images/disk_usage.go @@ -5,6 +5,7 @@ import ( "fmt" "os" "path/filepath" + "syscall" ) // totalReadyImageBytesFromMetadata sums ready image sizes directly from metadata.json files. @@ -14,6 +15,7 @@ import ( // files found in the digest directory so we do not undercount host disk usage. func totalReadyImageBytesFromMetadata(imagesDir string) (int64, error) { var total int64 + seenRootfs := make(map[rootfsIdentity]struct{}) err := filepath.Walk(imagesDir, func(path string, info os.FileInfo, err error) error { if err != nil { @@ -28,7 +30,7 @@ func totalReadyImageBytesFromMetadata(imagesDir string) (int64, error) { data, err := os.ReadFile(path) if err != nil { - rootfsBytes, fallbackErr := totalRootfsBytesInDigestDir(filepath.Dir(path)) + rootfsBytes, fallbackErr := totalUniqueRootfsBytesInDigestDir(filepath.Dir(path), seenRootfs) if fallbackErr == nil { total += rootfsBytes return nil @@ -38,18 +40,33 @@ func totalReadyImageBytesFromMetadata(imagesDir string) (int64, error) { var meta imageMetadata if err := json.Unmarshal(data, &meta); err != nil { - rootfsBytes, fallbackErr := totalRootfsBytesInDigestDir(filepath.Dir(path)) + rootfsBytes, fallbackErr := totalUniqueRootfsBytesInDigestDir(filepath.Dir(path), seenRootfs) if fallbackErr == nil { total += rootfsBytes return nil } return fmt.Errorf("unmarshal image metadata %s: %w", path, err) } - if meta.Status == StatusReady && meta.SizeBytes > 0 { - total += meta.SizeBytes - return nil - } if meta.Status == StatusReady { + rootfsPaths, globErr := filepath.Glob(filepath.Join(filepath.Dir(path), "rootfs.*")) + if globErr != nil { + return fmt.Errorf("find ready image rootfs for %s: %w", path, globErr) + } + for _, rootfsPath := range rootfsPaths { + rootfsInfo, statErr := os.Stat(rootfsPath) + if statErr != nil { + continue + } + if !markUniqueRootfs(rootfsInfo, seenRootfs) { + return nil + } + break + } + + if meta.SizeBytes > 0 { + total += meta.SizeBytes + return nil + } rootfsBytes, err := totalRootfsBytesInDigestDir(filepath.Dir(path)) if err != nil { return fmt.Errorf("stat ready image rootfs for %s: %w", path, err) @@ -172,3 +189,55 @@ func totalRootfsBytesInDigestDir(digestDir string) (int64, error) { } return total, nil } + +type rootfsIdentity struct { + dev uint64 + ino uint64 +} + +func markUniqueRootfs(info os.FileInfo, seen map[rootfsIdentity]struct{}) bool { + stat, ok := info.Sys().(*syscall.Stat_t) + if !ok { + return true + } + identity := rootfsIdentity{dev: uint64(stat.Dev), ino: uint64(stat.Ino)} + if _, exists := seen[identity]; exists { + return false + } + seen[identity] = struct{}{} + return true +} + +func totalUniqueRootfsBytesInDigestDir(digestDir string, seen map[rootfsIdentity]struct{}) (int64, error) { + rootfsPaths, err := filepath.Glob(filepath.Join(digestDir, "rootfs.*")) + if err != nil { + return 0, err + } + if len(rootfsPaths) == 0 { + return 0, os.ErrNotExist + } + + var total int64 + found := false + for _, rootfsPath := range rootfsPaths { + info, err := os.Stat(rootfsPath) + if err != nil { + if os.IsNotExist(err) { + continue + } + return 0, err + } + if info.IsDir() { + continue + } + found = true + if !markUniqueRootfs(info, seen) { + continue + } + total += info.Size() + } + if !found { + return 0, os.ErrNotExist + } + return total, nil +} diff --git a/lib/images/disk_usage_test.go b/lib/images/disk_usage_test.go index 36e5bc0df..0de352d48 100644 --- a/lib/images/disk_usage_test.go +++ b/lib/images/disk_usage_test.go @@ -9,8 +9,6 @@ import ( ) func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForMalformedMetadata(t *testing.T) { - t.Parallel() - imagesDir := t.TempDir() digestDir := filepath.Join(imagesDir, "docker.io", "library", "alpine", "sha256deadbeef") require.NoError(t, os.MkdirAll(digestDir, 0o755)) @@ -22,9 +20,45 @@ func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForMalformedMetadata require.Equal(t, int64(len("rootfs-data")), total) } -func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForReadyImageWithoutSize(t *testing.T) { - t.Parallel() +func TestTotalReadyImageBytesFromMetadata_DeduplicatesHardLinkedAliases(t *testing.T) { + imagesDir := t.TempDir() + sourceDir := filepath.Join(imagesDir, "source", "digest") + targetDir := filepath.Join(imagesDir, "target", "digest") + require.NoError(t, os.MkdirAll(sourceDir, 0o755)) + require.NoError(t, os.MkdirAll(targetDir, 0o755)) + + sourceRootfs := filepath.Join(sourceDir, "rootfs.erofs") + targetRootfs := filepath.Join(targetDir, "rootfs.erofs") + require.NoError(t, os.WriteFile(sourceRootfs, []byte("shared-rootfs"), 0o644)) + require.NoError(t, os.Link(sourceRootfs, targetRootfs)) + metadata := []byte(`{"status":"ready","size_bytes":13}`) + require.NoError(t, os.WriteFile(filepath.Join(sourceDir, "metadata.json"), metadata, 0o644)) + require.NoError(t, os.WriteFile(filepath.Join(targetDir, "metadata.json"), metadata, 0o644)) + + total, err := totalReadyImageBytesFromMetadata(imagesDir) + require.NoError(t, err) + require.Equal(t, int64(len("shared-rootfs")), total) +} +func TestTotalReadyImageBytesFromMetadata_DeduplicatesMalformedAliases(t *testing.T) { + imagesDir := t.TempDir() + malformedDir := filepath.Join(imagesDir, "a-malformed", "digest") + validDir := filepath.Join(imagesDir, "b-valid", "digest") + require.NoError(t, os.MkdirAll(malformedDir, 0o755)) + require.NoError(t, os.MkdirAll(validDir, 0o755)) + + rootfs := filepath.Join(malformedDir, "rootfs.erofs") + require.NoError(t, os.WriteFile(rootfs, []byte("shared-rootfs"), 0o644)) + require.NoError(t, os.Link(rootfs, filepath.Join(validDir, "rootfs.erofs"))) + require.NoError(t, os.WriteFile(filepath.Join(malformedDir, "metadata.json"), []byte("{not-json"), 0o644)) + require.NoError(t, os.WriteFile(filepath.Join(validDir, "metadata.json"), []byte(`{"status":"ready","size_bytes":13}`), 0o644)) + + total, err := totalReadyImageBytesFromMetadata(imagesDir) + require.NoError(t, err) + require.Equal(t, int64(len("shared-rootfs")), total) +} + +func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForReadyImageWithoutSize(t *testing.T) { imagesDir := t.TempDir() digestDir := filepath.Join(imagesDir, "docker.io", "library", "alpine", "sha256deadbeef") require.NoError(t, os.MkdirAll(digestDir, 0o755)) diff --git a/lib/images/manager.go b/lib/images/manager.go index c1a0da0ac..f89fb6678 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -72,6 +72,7 @@ type manager struct { queue *queue.Queue createMu sync.Mutex diskUsageMu sync.RWMutex + tagGenerations map[string]uint64 diskUsageLoaded bool readyImageBytes int64 ociCacheBytes int64 @@ -99,6 +100,7 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma inflightPulls: make(map[string]*inflightImagePull), borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, readySubscribers: make(map[string][]chan StatusEvent), + tagGenerations: make(map[string]uint64), } // Initialize metrics if meter is provided @@ -111,6 +113,7 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma } m.RecoverInterruptedBuilds() + promoteLegacyImages(m.paths) return m, nil } @@ -141,7 +144,7 @@ func (m *manager) ListImages(ctx context.Context) ([]Image, error) { images := make([]Image, 0, len(metas)) for _, meta := range metas { - images = append(images, *meta.toImage()) + images = append(images, *meta.toImageFor(meta.Name)) } return images, nil @@ -205,35 +208,9 @@ func (m *manager) CreateImage(ctx context.Context, req CreateImageRequest) (*Ima m.createMu.Lock() defer m.createMu.Unlock() - // Check if we already have this digest (deduplication) - if meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()); err == nil { - // Don't cache failed builds - allow retry - if meta.Status == StatusFailed { - // Clean up the failed build directory so we can retry - digestDir := filepath.Join(m.paths.ImagesDir(), ref.Repository(), ref.DigestHex()) - os.RemoveAll(digestDir) - // Fall through to re-queue the build - } else { - // We have this digest already (ready, pending, pulling, or converting). - // last-pull-wins: repoint the tag to it (see createTagSymlink). - if meta.Status == StatusReady { - if ref.Tag() != "" { - createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } - return meta.toImage(), nil - } - if !m.inflightCredentialsMatch(ref.Digest(), req.Credentials) { - return nil, fmt.Errorf("%w: retry after the current pull completes", ErrCredentialConflict) - } - img := meta.toImage() - if meta.Status == StatusPending { - img.QueuePosition = m.queue.GetPosition(meta.Digest) - } - return img, nil - } + if img, found, err := m.reuseExistingImage(ref, req.Credentials, req.Tags); found || err != nil { + return img, err } - - // Don't have this digest yet, queue the build return m.createAndQueueImage(ref, req, platform) } @@ -263,30 +240,113 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest m.createMu.Lock() defer m.createMu.Unlock() - // Check if we already have this digest (deduplication) - if meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()); err == nil { - // Don't cache failed builds - allow retry - if meta.Status == StatusFailed { - digestDir := filepath.Join(m.paths.ImagesDir(), ref.Repository(), ref.DigestHex()) - os.RemoveAll(digestDir) - // Fall through to re-queue the build - } else { - // We have this digest already; last-pull-wins repoints the tag. - if meta.Status == StatusReady && ref.Tag() != "" { - createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } - img := meta.toImage() - if meta.Status == StatusPending { - img.QueuePosition = m.queue.GetPosition(meta.Digest) - } - return img, nil - } + if img, found, err := m.reuseExistingImage(ref, nil, nil); found || err != nil { + return img, err } - // Don't have this digest yet, queue the build return m.createAndQueueImage(ref, CreateImageRequest{Name: imageRef}, hostPlatform()) } +func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthConfig, resourceTags tags.Tags) (*Image, bool, error) { + meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()) + if err != nil { + return nil, false, nil + } + if meta.Status == StatusFailed { + if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { + return nil, true, fmt.Errorf("remove failed image: %w", err) + } + return nil, false, nil + } + if err := m.updateExistingReference(meta, ref, resourceTags); err != nil { + return nil, true, fmt.Errorf("update image reference: %w", err) + } + img := meta.toImageFor(ref.String()) + if meta.Status == StatusReady { + return img, true, nil + } + if !m.inflightCredentialsMatch(ref.Digest(), credentials) { + return nil, true, fmt.Errorf("%w: retry after the current pull completes", ErrCredentialConflict) + } + if meta.Status == StatusPending { + img.QueuePosition = m.queue.GetPosition(meta.Digest) + } + return img, true, nil +} + +func (m *manager) updateExistingReference(meta *imageMetadata, ref *ResolvedRef, resourceTags tags.Tags) error { + if ref.Tag() == "" { + return nil + } + if meta.Status == StatusReady { + m.nextTagGeneration(ref.Repository(), ref.Tag()) + if err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()); err != nil { + return err + } + } else if err := m.recordPendingTag(meta, ref); err != nil { + return err + } + setReferenceTags(meta, ref.String(), resourceTags) + return writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta) +} + +func tagGenerationKey(repository, tag string) string { + return repository + ":" + tag +} + +func (m *manager) nextTagGeneration(repository, tag string) uint64 { + key := tagGenerationKey(repository, tag) + m.tagGenerations[key]++ + return m.tagGenerations[key] +} + +func (m *manager) recordPendingTag(meta *imageMetadata, ref *ResolvedRef) error { + if meta.RequestedTag == ref.Tag() && strings.HasPrefix(meta.Name, ref.Repository()+":") { + return ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } + for _, claim := range meta.TagClaims { + if claim.Repository == ref.Repository() && claim.Tag == ref.Tag() { + return ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } + } + previous, err := resolveTag(m.paths, ref.Repository(), ref.Tag()) + if err != nil && !errors.Is(err, ErrNotFound) { + return err + } + meta.TagClaims = append(meta.TagClaims, imageTagClaim{ + Repository: ref.Repository(), + Tag: ref.Tag(), + PreviousTagDigest: previous, + TagGeneration: m.nextTagGeneration(ref.Repository(), ref.Tag()), + }) + return ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) +} + +func (m *manager) restoreTagGenerations(metas []*imageMetadata) { + m.createMu.Lock() + defer m.createMu.Unlock() + for _, meta := range metas { + ref, err := ParseNormalizedRef(meta.Name) + if err != nil { + continue + } + m.restoreTagGeneration(ref.Repository(), meta.RequestedTag, meta.TagGeneration) + for _, claim := range meta.TagClaims { + m.restoreTagGeneration(claim.Repository, claim.Tag, claim.TagGeneration) + } + } +} + +func (m *manager) restoreTagGeneration(repository, tag string, generation uint64) { + if tag == "" { + return + } + key := tagGenerationKey(repository, tag) + if generation > m.tagGenerations[key] { + m.tagGenerations[key] = generation + } +} + func (m *manager) inflightCredentialsMatch(digest string, credentials *authn.AuthConfig) bool { inflight := m.inflightPulls[digest] var existingFingerprint [32]byte @@ -367,22 +427,43 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, // Record the requested platform optimistically while the build is pending; // buildImage overwrites it with the authoritative manifest platform once the // image config is pulled. + previousTagDigest := "" + var err error + if ref.Tag() != "" { + previousTagDigest, err = resolveTag(m.paths, ref.Repository(), ref.Tag()) + if err != nil && !errors.Is(err, ErrNotFound) { + return nil, fmt.Errorf("resolve existing image tag: %w", err) + } + } + tagGeneration := uint64(0) + if ref.Tag() != "" { + tagGeneration = m.nextTagGeneration(ref.Repository(), ref.Tag()) + } meta := &imageMetadata{ - Name: ref.String(), - Digest: ref.Digest(), - Platform: requestedPlatform.String(), - Status: StatusPending, - Request: &storedReq, - BorrowedAuth: req.Credentials != nil, - BuildID: uuid.New().String(), - Tags: tags.Clone(req.Tags), - CreatedAt: time.Now(), + Name: ref.String(), + Digest: ref.Digest(), + Platform: requestedPlatform.String(), + Status: StatusPending, + Request: &storedReq, + BorrowedAuth: req.Credentials != nil, + BuildID: uuid.New().String(), + Tags: tags.Clone(req.Tags), + RequestedTag: ref.Tag(), + PreviousTagDigest: previousTagDigest, + TagGeneration: tagGeneration, + References: map[string]tags.Tags{ref.String(): tags.Clone(req.Tags)}, + CreatedAt: time.Now(), } // Write initial metadata if err := writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta); err != nil { return nil, fmt.Errorf("write initial metadata: %w", err) } + if ref.Tag() != "" && previousTagDigest == "" { + if err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()); err != nil { + return nil, fmt.Errorf("create pending image tag: %w", err) + } + } // Keep borrowed credentials outside the queued closure so their lifetime is // bounded even when this job waits behind another pull. @@ -406,7 +487,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, m.buildImage(ctx, ref, credentials, buildID) }, m.releaseInflightPull(ref.Digest(), inflight)) - img := meta.toImage() + img := meta.toImageFor(ref.String()) if queuePos > 0 { img.QueuePosition = &queuePos } @@ -456,7 +537,13 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials if meta.Status == StatusReady { // Another build completed first; last-pull-wins repoints the tag. if ref.Tag() != "" { - createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + m.createMu.Lock() + m.nextTagGeneration(ref.Repository(), ref.Tag()) + err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + m.createMu.Unlock() + if err != nil { + slog.Warn("failed to claim ready image tag", "repository", ref.Repository(), "tag", ref.Tag(), "error", err) + } } buildStatus = "success" return @@ -465,10 +552,12 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials m.updateStatusByDigest(ref, StatusConverting, nil, buildID) - diskPath := digestPath(m.paths, ref.Repository(), ref.DigestHex()) + diskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk + diskTempPath := diskPath + ".tmp-" + buildID + defer os.Remove(diskTempPath) // Use default image format (erofs on Linux, ext4 on Darwin) convertStart := time.Now() - diskSize, err := ExportRootfs(tempDir, diskPath, DefaultImageFormat) + diskSize, err := ExportRootfs(tempDir, diskTempPath, DefaultImageFormat) m.recordImageBuildPhase(ctx, ref.Digest(), "filesystem_export", time.Since(convertStart), phaseStatus(err), "not_applicable") if err != nil { m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("convert to %s: %w", DefaultImageFormat, err), buildID) @@ -476,7 +565,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } finalizeStart := time.Now() - err = m.finalizeImage(ref, result, diskSize, buildID) + err = m.finalizeImage(ref, result, diskSize, buildID, diskTempPath) m.recordImageBuildPhase(ctx, ref.Digest(), "finalize", time.Since(finalizeStart), phaseStatus(err), "not_applicable") if err != nil { if errors.Is(err, errStaleBuild) { @@ -489,7 +578,11 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials buildStatus = "success" } -func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID string) error { +func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string) error { + if diskTempPath != "" { + defer os.Remove(diskTempPath) + } + m.createMu.Lock() defer m.createMu.Unlock() @@ -499,6 +592,13 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i return errStaleBuild } + finalDiskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk + if err := installAtomically(finalDiskPath, func(path string) error { + return os.Rename(diskTempPath, path) + }); err != nil { + return fmt.Errorf("install image disk: %w", err) + } + // The pulled image config is the source of truth for the platform. var requestedPlatform string if meta.Request != nil { @@ -524,15 +624,39 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i } m.notifyReady(ref.DigestHex(), StatusReady, nil) - if ref.Tag() != "" { - if err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()); err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) - } - } + m.claimImageTags(ref, meta) m.refreshDiskUsageTotals() return nil } +func (m *manager) claimImageTags(ref *ResolvedRef, meta *imageMetadata) { + requestedTag := meta.RequestedTag + if requestedTag == "" { + requestedTag = ref.Tag() + } + m.claimTag(ref.Repository(), ref.DigestHex(), requestedTag, meta.PreviousTagDigest, meta.TagGeneration, meta.RequestedTag == "") + for _, claim := range meta.TagClaims { + m.claimTag(claim.Repository, ref.DigestHex(), claim.Tag, claim.PreviousTagDigest, claim.TagGeneration, false) + } +} + +func (m *manager) claimTag(repository, digestHex, tag, previous string, generation uint64, allowMissing bool) { + if tag == "" || m.tagGenerations[tagGenerationKey(repository, tag)] != generation { + return + } + current, err := resolveTag(m.paths, repository, tag) + if err != nil { + if !allowMissing || !errors.Is(err, ErrNotFound) { + return + } + } else if current != digestHex && current != previous { + return + } + if err := createTagSymlink(m.paths, repository, tag, digestHex); err != nil { + fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) + } +} + func phaseStatus(err error) string { if err != nil { return "failed" @@ -602,12 +726,14 @@ func (m *manager) RecoverInterruptedBuilds() { if err != nil { return // Best effort } + m.restoreTagGenerations(metas) // Sort by created_at to maintain FIFO order sort.Slice(metas, func(i, j int) bool { return metas[i].CreatedAt.Before(metas[j].CreatedAt) }) + seenDigests := make(map[string]struct{}) for _, meta := range metas { if meta.Status != StatusPending && meta.Status != StatusPulling && meta.Status != StatusConverting { continue @@ -615,6 +741,10 @@ func (m *manager) RecoverInterruptedBuilds() { if meta.Request == nil || meta.Digest == "" { continue } + if _, seen := seenDigests[meta.Digest]; seen { + continue + } + seenDigests[meta.Digest] = struct{}{} normalized, err := ParseNormalizedRef(meta.Name) if err != nil { continue @@ -659,7 +789,7 @@ func (m *manager) GetImage(ctx context.Context, name string) (*Image, error) { return nil, err } - img := meta.toImage() + img := meta.toImageFor(ref.String()) if meta.Status == StatusPending { img.QueuePosition = m.queue.GetPosition(meta.Digest) @@ -690,7 +820,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := deleteTagsForDigest(m.paths, repository, digestHex); err != nil { return err } - if err := deleteDigest(m.paths, repository, digestHex); err != nil { + if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, false); err != nil { return err } m.refreshDiskUsageTotals() @@ -698,6 +828,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { } tag := ref.Tag() + m.nextTagGeneration(repository, tag) // Resolve the tag to get the digest before deleting digestHex, err := resolveTag(m.paths, repository, tag) @@ -713,15 +844,12 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { // Check if the digest is now orphaned (no other tags reference it) count, err := countTagsForDigest(m.paths, repository, digestHex) if err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to count tags for digest %s: %v\n", digestHex, err) - return nil + return fmt.Errorf("count tags for digest %s: %w", digestHex, err) } if count == 0 { - // Digest is orphaned, delete it - if err := deleteDigest(m.paths, repository, digestHex); err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to delete orphaned digest %s: %v\n", digestHex, err) - return nil + if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, true); err != nil { + return fmt.Errorf("delete orphaned digest %s: %w", digestHex, err) } m.refreshDiskUsageTotals() } @@ -747,42 +875,45 @@ func (m *manager) TotalOCICacheBytes(ctx context.Context) (int64, error) { return ociCacheBytes, nil } -// WaitForReady blocks until the image reaches a terminal state (ready or failed) -// or the context is cancelled. -// -// The image may not exist yet when this is called (e.g., the registry's -// triggerConversion goroutine hasn't called ImportLocalImage yet), so we -// poll briefly for the image to appear before subscribing for notifications. -func (m *manager) WaitForReady(ctx context.Context, name string) error { - // Wait for the image to appear in the store. In the build flow, the - // registry triggers ImportLocalImage asynchronously after a push, so the - // image may not exist when the build manager calls WaitForReady. - const maxWaitForExist = 30 * time.Second - const pollInterval = 100 * time.Millisecond - - var img *Image - deadline := time.Now().Add(maxWaitForExist) - for { - got, err := m.GetImage(ctx, name) - if err == nil { - img = got - break +func (m *manager) findRequestedTagImage(ref *NormalizedRef) *Image { + metas, err := listAllMetadata(m.paths) + if err != nil { + return nil + } + var newest *imageMetadata + for _, meta := range metas { + requestedTag := meta.RequestedTag + if requestedTag == "" { + requestedTag = strings.TrimPrefix(meta.Name, ref.Repository()+":") } - if time.Now().After(deadline) { - return fmt.Errorf("get image: %w", err) + if requestedTag != ref.Tag() || !strings.HasPrefix(meta.Name, ref.Repository()+":") { + continue } - select { - case <-ctx.Done(): - return ctx.Err() - case <-time.After(pollInterval): + if newest == nil || newest.CreatedAt.Before(meta.CreatedAt) { + newest = meta } } + if newest == nil { + return nil + } + return newest.toImageFor(ref.String()) +} - // Check if already in terminal state - switch img.Status { - case StatusReady: +// WaitForReady blocks until the image reaches a terminal state (ready or failed) +// or the context is cancelled. +func (m *manager) WaitForReady(ctx context.Context, name string) error { + ref, err := ParseNormalizedRef(name) + if err != nil { + return fmt.Errorf("parse image name: %w", err) + } + img, err := m.waitForImage(ctx, name, ref) + if err != nil { + return err + } + if img.Status == StatusReady { return nil - case StatusFailed: + } + if img.Status == StatusFailed { return conversionFailedErr(img.Error, nil) } @@ -794,12 +925,12 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { defer m.unsubscribeFromReady(digestHex, ch) // Re-check after subscribing to close the race window - img, err := m.GetImage(ctx, name) + img, err = m.GetImage(ctx, ref.Repository()+"@"+img.Digest) if err == nil { - switch img.Status { - case StatusReady: + if img.Status == StatusReady { return nil - case StatusFailed: + } + if img.Status == StatusFailed { return conversionFailedErr(img.Error, nil) } } @@ -816,6 +947,33 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { } } +func (m *manager) waitForImage(ctx context.Context, name string, ref *NormalizedRef) (*Image, error) { + const maxWaitForExist = 30 * time.Second + const pollInterval = 100 * time.Millisecond + var lastErr error + deadline := time.Now().Add(maxWaitForExist) + for { + var img *Image + if !ref.IsDigest() { + img = m.findRequestedTagImage(ref) + } + if img == nil { + img, lastErr = m.GetImage(ctx, name) + } + if img != nil { + return img, nil + } + if time.Now().After(deadline) { + return nil, fmt.Errorf("get image: %w", lastErr) + } + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(pollInterval): + } + } +} + func conversionFailedErr(message *string, cause error) error { if cause != nil { return fmt.Errorf("image conversion failed: %w", cause) diff --git a/lib/images/manager_test.go b/lib/images/manager_test.go index 89d3fb4fe..8abd90a65 100644 --- a/lib/images/manager_test.go +++ b/lib/images/manager_test.go @@ -98,10 +98,12 @@ func TestCreateImage(t *testing.T) { require.NoError(t, err) require.NotEqual(t, 0, linkStat.Mode()&os.ModeSymlink, "should be a symlink") - // Verify symlink points to digest directory + // Verify symlink resolves to the digest linkTarget, err := os.Readlink(linkPath) require.NoError(t, err) - require.Equal(t, digestHex, linkTarget, "symlink should point to digest") + resolvedDigest, err := resolveTag(paths.New(dataDir), ref.Repository(), ref.Tag()) + require.NoError(t, err) + require.Equal(t, digestHex, resolvedDigest, "symlink should resolve to digest") t.Logf("Tag symlink: %s -> %s", linkPath, linkTarget) } @@ -395,6 +397,39 @@ func TestDeleteImagePreservesSharedDigest(t *testing.T) { require.True(t, os.IsNotExist(err), "digest directory should be deleted when last tag is removed") } +func TestDeleteImagePreservesCrossRepositoryContent(t *testing.T) { + p := paths.New(t.TempDir()) + mgr, err := NewManager(p, 1, nil) + require.NoError(t, err) + + const digest = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + meta := &imageMetadata{ + Name: "docker.io/library/alpine@sha256:" + digest, + Digest: "sha256:" + digest, + Status: StatusReady, + SizeBytes: 6, + CreatedAt: time.Now().UTC(), + } + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("rootfs"), 0o644)) + require.NoError(t, createTagSymlink(p, "docker.io/library/alpine", "latest", digest)) + require.NoError(t, createTagSymlink(p, "registry.example.com/app", "v1", digest)) + + require.NoError(t, mgr.DeleteImage(context.Background(), "docker.io/library/alpine:latest")) + _, err = mgr.GetImage(context.Background(), "registry.example.com/app:v1") + require.NoError(t, err) + _, err = os.Stat(p.ImageContentPath(digest)) + require.NoError(t, err) + + require.NoError(t, mgr.DeleteImage(context.Background(), "registry.example.com/app:v1")) + _, err = os.Stat(p.ImageContentDir(digest)) + require.NoError(t, err) + + require.NoError(t, mgr.DeleteImage(context.Background(), "docker.io/library/alpine@sha256:"+digest)) + _, err = os.Stat(p.ImageContentDir(digest)) + require.ErrorIs(t, err, os.ErrNotExist) +} + func TestNormalizedRefParsing(t *testing.T) { tests := []struct { input string @@ -671,6 +706,9 @@ func TestDeleteAndRecreateDuringBuildTail(t *testing.T) { require.NoError(t, err) require.NotEqual(t, firstMeta.BuildID, currentMeta.BuildID) require.NoError(t, m.DeleteImage(ctx, repo+"@"+digestStr)) + // Deleting a digest whose build is still in flight keeps the shared content + // so the running build can finish; re-import joins that build instead of + // starting a second one for the same digest. recreatedAgain, err := m.ImportLocalImage(ctx, repo, tag, digestStr) require.NoError(t, err) require.Equal(t, StatusPending, recreatedAgain.Status) @@ -678,7 +716,7 @@ func TestDeleteAndRecreateDuringBuildTail(t *testing.T) { require.Equal(t, 1, *recreatedAgain.QueuePosition) latestMeta, err := readMetadata(p, repo, digestHex) require.NoError(t, err) - require.NotEqual(t, currentMeta.BuildID, latestMeta.BuildID) + require.Equal(t, currentMeta.BuildID, latestMeta.BuildID) currentMeta = latestMeta waitCtx, cancelWait := context.WithCancel(ctx) @@ -698,7 +736,7 @@ func TestDeleteAndRecreateDuringBuildTail(t *testing.T) { m.updateStatusByDigest(staleRef, StatusFailed, errors.New("stale build"), firstMeta.BuildID) staleResult, _, _, err := m.ociClient.extractOCIImageDetails(digestHex) require.NoError(t, err) - require.ErrorIs(t, m.finalizeImage(staleRef, &pullResult{Metadata: staleResult}, 1, firstMeta.BuildID), errStaleBuild) + require.ErrorIs(t, m.finalizeImage(staleRef, &pullResult{Metadata: staleResult}, 1, firstMeta.BuildID, ""), errStaleBuild) currentMeta, err = readMetadata(p, repo, digestHex) require.NoError(t, err) require.Equal(t, StatusPending, currentMeta.Status) diff --git a/lib/images/storage.go b/lib/images/storage.go index 8764400c1..2e3430851 100644 --- a/lib/images/storage.go +++ b/lib/images/storage.go @@ -14,22 +14,55 @@ import ( ) type imageMetadata struct { - Name string `json:"name"` // Normalized ref (tag or digest) - Digest string `json:"digest"` // Always present: sha256:... - Platform string `json:"platform,omitempty"` - Status string `json:"status"` - Error *string `json:"error,omitempty"` - Request *CreateImageRequest `json:"request,omitempty"` - SizeBytes int64 `json:"size_bytes"` - Entrypoint []string `json:"entrypoint,omitempty"` - Cmd []string `json:"cmd,omitempty"` - Env map[string]string `json:"env,omitempty"` - Labels map[string]string `json:"labels,omitempty"` - Tags tags.Tags `json:"tags,omitempty"` - WorkingDir string `json:"working_dir,omitempty"` - CreatedAt time.Time `json:"created_at"` - BorrowedAuth bool `json:"borrowed_auth,omitempty"` - BuildID string `json:"build_id,omitempty"` + Name string `json:"name"` // Normalized ref (tag or digest) + Digest string `json:"digest"` // Always present: sha256:... + Platform string `json:"platform,omitempty"` + Status string `json:"status"` + Error *string `json:"error,omitempty"` + Request *CreateImageRequest `json:"request,omitempty"` + SizeBytes int64 `json:"size_bytes"` + Entrypoint []string `json:"entrypoint,omitempty"` + Cmd []string `json:"cmd,omitempty"` + Env map[string]string `json:"env,omitempty"` + Labels map[string]string `json:"labels,omitempty"` + Tags tags.Tags `json:"tags,omitempty"` + WorkingDir string `json:"working_dir,omitempty"` + CreatedAt time.Time `json:"created_at"` + BorrowedAuth bool `json:"borrowed_auth,omitempty"` + BuildID string `json:"build_id,omitempty"` + RequestedTag string `json:"requested_tag,omitempty"` + PreviousTagDigest string `json:"previous_tag_digest,omitempty"` + TagGeneration uint64 `json:"tag_generation,omitempty"` + References map[string]tags.Tags `json:"references,omitempty"` + TagClaims []imageTagClaim `json:"tag_claims,omitempty"` +} + +type imageTagClaim struct { + Repository string `json:"repository"` + Tag string `json:"tag"` + PreviousTagDigest string `json:"previous_tag_digest,omitempty"` + TagGeneration uint64 `json:"tag_generation,omitempty"` +} + +func referenceTags(meta *imageMetadata, reference string) (tags.Tags, bool) { + resourceTags, ok := meta.References[reference] + return tags.Clone(resourceTags), ok +} + +func setReferenceTags(meta *imageMetadata, reference string, resourceTags tags.Tags) { + if meta.References == nil { + meta.References = make(map[string]tags.Tags) + } + meta.References[reference] = tags.Clone(resourceTags) +} + +func (m *imageMetadata) toImageFor(reference string) *Image { + img := m.toImage() + if resourceTags, ok := referenceTags(m, reference); ok { + img.Tags = resourceTags + } + img.Name = reference + return img } func (m *imageMetadata) toImage() *Image { @@ -77,16 +110,75 @@ func (m *imageMetadata) toImage() *Image { return img } -// digestDir returns the directory for a specific digest -// e.g., /var/lib/hypeman/images/docker.io/library/alpine/abc123def456... +// imageLayout contains the paths for one image layout. +type imageLayout struct { + dir string + metadata string + disk string + content bool +} + +// resolveImageLayout selects one layout for all operations on a digest. A +// complete legacy image remains authoritative while content is incomplete; +// once content metadata is ready, content becomes canonical even if its disk +// is missing so callers report the corruption instead of mixing layouts. +func resolveImageLayout(p *paths.Paths, repository, digestHex string) imageLayout { + legacy := imageLayout{ + dir: p.ImageDigestDir(repository, digestHex), + metadata: p.ImageMetadata(repository, digestHex), + disk: p.ImageDigestPath(repository, digestHex), + } + content := contentLayout(p, digestHex) + + if legacyImageExists(p, repository, digestHex) { + contentStatus, contentOK := metadataStatus(content.metadata) + if !contentOK || contentStatus != StatusReady { + return legacy + } + } + if pathExists(content.metadata) || pathExists(content.disk) { + return content + } + if pathExists(legacy.metadata) { + return legacy + } + return content +} + +func contentLayout(p *paths.Paths, digestHex string) imageLayout { + return imageLayout{ + dir: p.ImageContentDir(digestHex), + metadata: p.ImageContentMetadata(digestHex), + disk: p.ImageContentPath(digestHex), + content: true, + } +} + +func pathExists(path string) bool { + _, err := os.Stat(path) + return err == nil +} + func digestDir(p *paths.Paths, repository, digestHex string) string { - return p.ImageDigestDir(repository, digestHex) + return resolveImageLayout(p, repository, digestHex).dir } -// digestPath returns the path to the rootfs disk file for a digest -// Uses .erofs on Linux (compressed) or .ext4 on Darwin (VZ kernel lacks erofs support) -func digestPath(p *paths.Paths, repository, digestHex string) string { - return p.ImageDigestPath(repository, digestHex) +func legacyImageExists(p *paths.Paths, repository, digestHex string) bool { + _, metadataErr := os.Stat(p.ImageMetadata(repository, digestHex)) + _, diskErr := os.Stat(p.ImageDigestPath(repository, digestHex)) + return metadataErr == nil && diskErr == nil +} + +func metadataStatus(path string) (string, bool) { + data, err := os.ReadFile(path) + if err != nil { + return "", false + } + var meta imageMetadata + if err := json.Unmarshal(data, &meta); err != nil { + return "", false + } + return meta.Status, true } // GetDiskPath returns the filesystem path to an image's rootfs disk file (public for instances manager) @@ -100,25 +192,32 @@ func GetDiskPath(p *paths.Paths, imageName string, digest string) (string, error // Extract digest hex (remove "sha256:" prefix) digestHex := strings.TrimPrefix(digest, "sha256:") - return digestPath(p, ref.Repository(), digestHex), nil + return resolveImageLayout(p, ref.Repository(), digestHex).disk, nil +} + +func digestPath(p *paths.Paths, repository, digestHex string) string { + return resolveImageLayout(p, repository, digestHex).disk } -// metadataPath returns the path to metadata.json for a digest func metadataPath(p *paths.Paths, repository, digestHex string) string { - return p.ImageMetadata(repository, digestHex) + return resolveImageLayout(p, repository, digestHex).metadata } -// tagSymlinkPath returns the path to a tag symlink -// e.g., /var/lib/hypeman/images/docker.io/library/alpine/latest func tagSymlinkPath(p *paths.Paths, repository, tag string) string { + newPath := p.ImageRepositoryTagSymlink(repository, tag) + if _, err := os.Lstat(newPath); err == nil { + return newPath + } return p.ImageTagSymlink(repository, tag) } -// writeMetadata writes metadata for a digest func writeMetadata(p *paths.Paths, repository, digestHex string, meta *imageMetadata) error { - dir := digestDir(p, repository, digestHex) - if err := os.MkdirAll(dir, 0755); err != nil { - return fmt.Errorf("create digest directory: %w", err) + return writeMetadataFile(resolveImageLayout(p, repository, digestHex).metadata, meta) +} + +func writeMetadataFile(path string, meta *imageMetadata) error { + if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil { + return fmt.Errorf("create metadata directory: %w", err) } data, err := json.MarshalIndent(meta, "", " ") @@ -126,23 +225,27 @@ func writeMetadata(p *paths.Paths, repository, digestHex string, meta *imageMeta return fmt.Errorf("marshal metadata: %w", err) } - tempPath := metadataPath(p, repository, digestHex) + ".tmp" + tempPath := path + ".tmp" if err := os.WriteFile(tempPath, data, 0644); err != nil { return fmt.Errorf("write temp metadata: %w", err) } - - finalPath := metadataPath(p, repository, digestHex) - if err := os.Rename(tempPath, finalPath); err != nil { - os.Remove(tempPath) + if err := os.Rename(tempPath, path); err != nil { + _ = os.Remove(tempPath) return fmt.Errorf("rename metadata: %w", err) } - return nil } -// readMetadata reads metadata for a digest func readMetadata(p *paths.Paths, repository, digestHex string) (*imageMetadata, error) { - path := metadataPath(p, repository, digestHex) + return readMetadataAt(resolveImageLayout(p, repository, digestHex)) +} + +func readContentMetadata(p *paths.Paths, digestHex string) (*imageMetadata, error) { + return readMetadataAt(contentLayout(p, digestHex)) +} + +func readMetadataAt(layout imageLayout) (*imageMetadata, error) { + path := layout.metadata data, err := os.ReadFile(path) if err != nil { if os.IsNotExist(err) { @@ -157,10 +260,9 @@ func readMetadata(p *paths.Paths, repository, digestHex string) (*imageMetadata, } if meta.Status == StatusReady { - diskPath := digestPath(p, repository, digestHex) - if _, err := os.Stat(diskPath); err != nil { + if _, err := os.Stat(layout.disk); err != nil { if os.IsNotExist(err) { - return nil, fmt.Errorf("disk image missing: %s", diskPath) + return nil, fmt.Errorf("disk image missing: %s", layout.disk) } return nil, fmt.Errorf("stat disk image: %w", err) } @@ -169,220 +271,231 @@ func readMetadata(p *paths.Paths, repository, digestHex string) (*imageMetadata, return &meta, nil } -// createTagSymlink creates or updates a tag symlink to point to a digest (only -// if the digest dir exists and the build is ready). -// -// Tag ownership is Docker last-pull-wins: the most recent pull of a tag always -// owns the symlink, regardless of platform. An earlier gate only repointed for -// host-native pulls, which silently stranded emulated variants (e.g. -// `pull --platform linux/amd64 alpine:3.19` could never make `image get` report -// amd64) and was non-recoverable. Always repointing is symmetric and matches -// Docker; callers repoint unconditionally on a ready digest. -func createTagSymlink(p *paths.Paths, repository, tag, digestHex string) error { - linkPath := tagSymlinkPath(p, repository, tag) - targetPath := digestHex // Relative path (just the digest hex) - - // Ensure parent directory exists - if err := os.MkdirAll(filepath.Dir(linkPath), 0755); err != nil { - return fmt.Errorf("create parent directory: %w", err) +func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata) error { + contentReady := false + if contentMeta, err := readContentMetadata(p, digestHex); err == nil { + contentReady = contentMeta.Status == StatusReady } - - // Remove existing symlink if present - os.Remove(linkPath) - - // Create new symlink - if err := os.Symlink(targetPath, linkPath); err != nil { - return fmt.Errorf("create symlink: %w", err) + if !contentReady { + sourceLayout := resolveImageLayout(p, sourceRepository, digestHex) + sourceDiskPath := sourceLayout.disk + if _, err := os.Stat(sourceDiskPath); err != nil { + return fmt.Errorf("stat source disk: %w", err) + } + if err := os.MkdirAll(p.ImageContentDir(digestHex), 0755); err != nil { + return fmt.Errorf("create content directory: %w", err) + } + contentMeta := *sourceMeta + contentMeta.Status = StatusConverting + if err := writeMetadataFile(p.ImageContentMetadata(digestHex), &contentMeta); err != nil { + return fmt.Errorf("write content metadata: %w", err) + } + if err := installAtomically(p.ImageContentPath(digestHex), func(path string) error { + return os.Link(sourceDiskPath, path) + }); err != nil { + return fmt.Errorf("link source disk: %w", err) + } + if err := writeMetadataFile(p.ImageContentMetadata(digestHex), sourceMeta); err != nil { + return fmt.Errorf("finalize content metadata: %w", err) + } } - return nil + return promoteLegacyTags(p, sourceRepository, digestHex) } -// resolveTag follows a tag symlink to get the digest hex -func resolveTag(p *paths.Paths, repository, tag string) (string, error) { - linkPath := tagSymlinkPath(p, repository, tag) - - // Read the symlink - target, err := os.Readlink(linkPath) +func promoteLegacyTags(p *paths.Paths, repository, digestHex string) error { + if _, err := os.Stat(p.ImageDigestDir(repository, digestHex)); err != nil { + return nil + } + tags, err := listTags(p, repository) if err != nil { - if os.IsNotExist(err) { - return "", ErrNotFound + return err + } + staged := make([]stagedTagSymlink, 0, len(tags)) + defer func() { + for _, ref := range staged { + _ = os.RemoveAll(ref.tempDir) } - return "", fmt.Errorf("read symlink: %w", err) + }() + for _, tag := range tags { + target, err := resolveTag(p, repository, tag) + if err != nil || target != digestHex { + continue + } + ref, err := stageTagSymlink(p, repository, tag, digestHex) + if err != nil { + return fmt.Errorf("stage legacy tag %s: %w", tag, err) + } + staged = append(staged, ref) } - - // Validate it's just a digest hex (not an absolute path) - if filepath.IsAbs(target) || strings.Contains(target, "/") { - return "", fmt.Errorf("invalid symlink target: %s", target) + for i, ref := range staged { + if err := os.Rename(ref.tempPath, ref.linkPath); err != nil { + if rollbackErr := rollbackTagSymlinks(staged[:i]); rollbackErr != nil { + return errors.Join(fmt.Errorf("promote legacy tag: %w", err), rollbackErr) + } + return fmt.Errorf("promote legacy tag: %w", err) + } } - - return target, nil + return nil } -// listTags returns all tags for a repository -func listTags(p *paths.Paths, repository string) ([]string, error) { - repoDir := p.ImageRepositoryDir(repository) - - entries, err := os.ReadDir(repoDir) +func installAtomically(path string, install func(string) error) error { + if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil { + return err + } + tempDir, err := os.MkdirTemp(filepath.Dir(path), ".install-*") if err != nil { - if os.IsNotExist(err) { - return nil, nil - } - return nil, fmt.Errorf("read repository directory: %w", err) + return err } - - var tags []string - for _, entry := range entries { - // Check if it's a symlink - info, err := os.Lstat(filepath.Join(repoDir, entry.Name())) - if err != nil { - continue - } - - if info.Mode()&os.ModeSymlink != 0 { - tags = append(tags, entry.Name()) - } + defer os.RemoveAll(tempDir) + tempPath := filepath.Join(tempDir, filepath.Base(path)) + if err := install(tempPath); err != nil { + return err } + return os.Rename(tempPath, path) +} - return tags, nil +type symlinkState struct { + exists bool + target string } -// listAllMetadata returns one metadata record per digest across all repositories. -// Tagged images are discovered through tag symlinks, and digest-only images are -// discovered directly from their metadata.json files. -func listAllMetadata(p *paths.Paths) ([]*imageMetadata, error) { - imagesDir := p.ImagesDir() - seen := make(map[string]struct{}) - metas := make([]*imageMetadata, 0) +type stagedTagSymlink struct { + repository string + tag string + linkPath string + tempDir string + tempPath string + previous symlinkState +} - err := filepath.Walk(imagesDir, func(path string, info os.FileInfo, err error) error { +func stageTagSymlink(p *paths.Paths, repository, tag, digestHex string) (stagedTagSymlink, error) { + layout := resolveImageLayout(p, repository, digestHex) + linkPath := p.ImageTagSymlink(repository, tag) + targetPath := digestHex // Relative path (just the digest hex) + if layout.content { + linkPath = p.ImageRepositoryTagSymlink(repository, tag) + var err error + targetPath, err = filepath.Rel(filepath.Dir(linkPath), layout.dir) if err != nil { - return nil // Skip errors + return stagedTagSymlink{}, fmt.Errorf("calculate content symlink target: %w", err) } + } + previous, err := readSymlinkState(linkPath) + if err != nil { + return stagedTagSymlink{}, err + } + if err := os.MkdirAll(filepath.Dir(linkPath), 0755); err != nil { + return stagedTagSymlink{}, fmt.Errorf("create parent directory: %w", err) + } + tempDir, err := os.MkdirTemp(filepath.Dir(linkPath), ".tag-stage-*") + if err != nil { + return stagedTagSymlink{}, fmt.Errorf("create temporary tag directory: %w", err) + } + tempPath := filepath.Join(tempDir, filepath.Base(linkPath)) + if err := os.Symlink(targetPath, tempPath); err != nil { + _ = os.RemoveAll(tempDir) + return stagedTagSymlink{}, fmt.Errorf("create temporary tag symlink: %w", err) + } + return stagedTagSymlink{ + repository: repository, + tag: tag, + linkPath: linkPath, + tempDir: tempDir, + tempPath: tempPath, + previous: previous, + }, nil +} - switch { - case info.Mode()&os.ModeSymlink != 0: - digestHex, err := os.Readlink(path) - if err != nil { - return nil // Skip invalid symlinks - } - - repository, err := filepath.Rel(imagesDir, filepath.Dir(path)) - if err != nil { - return nil - } - - return appendMetadataIfNew(p, repository, digestHex, seen, &metas) - case !info.IsDir() && info.Name() == "metadata.json": - digestHex := filepath.Base(filepath.Dir(path)) - repository, err := filepath.Rel(imagesDir, filepath.Dir(filepath.Dir(path))) - if err != nil { - return nil - } - - return appendMetadataIfNew(p, repository, digestHex, seen, &metas) - default: - return nil +func readSymlinkState(path string) (symlinkState, error) { + target, err := os.Readlink(path) + if err != nil { + if os.IsNotExist(err) { + return symlinkState{}, nil } - }) - - if err != nil && !os.IsNotExist(err) { - return nil, fmt.Errorf("walk images directory: %w", err) + return symlinkState{}, fmt.Errorf("read existing tag symlink: %w", err) } - - return metas, nil + return symlinkState{exists: true, target: target}, nil } -func appendMetadataIfNew(p *paths.Paths, repository, digestHex string, seen map[string]struct{}, metas *[]*imageMetadata) error { - key := repository + "@" + digestHex - if _, ok := seen[key]; ok { +func restoreSymlinkState(path string, state symlinkState) error { + if !state.exists { + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("remove tag symlink during rollback: %w", err) + } return nil } - - meta, err := readMetadata(p, repository, digestHex) - if err != nil { - return nil // Skip if metadata can't be read + if err := installAtomically(path, func(tempPath string) error { + return os.Symlink(state.target, tempPath) + }); err != nil { + return fmt.Errorf("restore tag symlink during rollback: %w", err) } - - seen[key] = struct{}{} - *metas = append(*metas, meta) return nil } -// digestExists checks if a digest directory exists -func digestExists(p *paths.Paths, repository, digestHex string) bool { - dir := digestDir(p, repository, digestHex) - _, err := os.Stat(dir) - return err == nil -} - -// deleteTag removes a tag symlink (does not delete the digest directory) -func deleteTag(p *paths.Paths, repository, tag string) error { - linkPath := tagSymlinkPath(p, repository, tag) - - // Check if symlink exists - if _, err := os.Lstat(linkPath); err != nil { - if os.IsNotExist(err) { - return ErrNotFound +func rollbackTagSymlinks(refs []stagedTagSymlink) error { + var rollbackErrs []error + for i := len(refs) - 1; i >= 0; i-- { + if err := restoreSymlinkState(refs[i].linkPath, refs[i].previous); err != nil { + rollbackErrs = append(rollbackErrs, err) } - return fmt.Errorf("stat symlink: %w", err) } + return errors.Join(rollbackErrs...) +} - // Remove symlink - if err := os.Remove(linkPath); err != nil { - return fmt.Errorf("remove symlink: %w", err) +func removeStaleTagSymlink(p *paths.Paths, ref *stagedTagSymlink) error { + stalePath := p.ImageTagSymlink(ref.repository, ref.tag) + if stalePath == ref.linkPath { + stalePath = p.ImageRepositoryTagSymlink(ref.repository, ref.tag) + } + if err := os.Remove(stalePath); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("remove stale tag symlink: %w", err) } - return nil } -// countTagsForDigest counts how many tags in a repository point to a given digest -func countTagsForDigest(p *paths.Paths, repository, digestHex string) (int, error) { - tags, err := listTags(p, repository) +func createTagSymlink(p *paths.Paths, repository, tag, digestHex string) error { + ref, err := stageTagSymlink(p, repository, tag, digestHex) if err != nil { - return 0, err + return fmt.Errorf("stage tag symlink: %w", err) } - - count := 0 - for _, tag := range tags { - target, err := resolveTag(p, repository, tag) - if err != nil { - continue - } - if target == digestHex { - count++ - } + if err := os.Rename(ref.tempPath, ref.linkPath); err != nil { + _ = os.RemoveAll(ref.tempDir) + return fmt.Errorf("install tag symlink: %w", err) + } + if err := removeStaleTagSymlink(p, &ref); err != nil { + fmt.Fprintf(os.Stderr, "Warning: failed to remove stale tag symlink %s: %v\n", tag, err) } - return count, nil + _ = os.RemoveAll(ref.tempDir) + return nil } -func deleteTagsForDigest(p *paths.Paths, repository, digestHex string) error { - tags, err := listTags(p, repository) +func resolveTag(p *paths.Paths, repository, tag string) (string, error) { + linkPath := tagSymlinkPath(p, repository, tag) + + target, err := os.Readlink(linkPath) if err != nil { - return err + if os.IsNotExist(err) { + return "", ErrNotFound + } + return "", fmt.Errorf("read symlink: %w", err) } - for _, tag := range tags { - target, err := resolveTag(p, repository, tag) - if err != nil { - continue - } - if target != digestHex { - continue - } - if err := deleteTag(p, repository, tag); err != nil && !errors.Is(err, ErrNotFound) { - return err + // Legacy links contain only the digest. New links point relatively into the + // shared content directory; validate that they resolve to that digest only. + if filepath.IsAbs(target) { + return "", fmt.Errorf("invalid symlink target: %s", target) + } + digestHex := filepath.Base(target) + if digestHex == "." || digestHex == string(filepath.Separator) { + return "", fmt.Errorf("invalid symlink target: %s", target) + } + if target != digestHex { + resolved := filepath.Clean(filepath.Join(filepath.Dir(linkPath), target)) + if resolved != filepath.Clean(p.ImageContentDir(digestHex)) { + return "", fmt.Errorf("invalid symlink target: %s", target) } } - return nil -} - -// deleteDigest removes a digest directory and all its contents -func deleteDigest(p *paths.Paths, repository, digestHex string) error { - dir := digestDir(p, repository, digestHex) - if err := os.RemoveAll(dir); err != nil { - return fmt.Errorf("remove digest directory: %w", err) - } - return nil + return digestHex, nil } diff --git a/lib/images/storage_refs.go b/lib/images/storage_refs.go new file mode 100644 index 000000000..ac1864c2f --- /dev/null +++ b/lib/images/storage_refs.go @@ -0,0 +1,357 @@ +package images + +import ( + "errors" + "fmt" + "os" + "path/filepath" + "strings" + + "github.com/kernel/hypeman/lib/paths" +) + +func ensurePendingTag(p *paths.Paths, repository, tag, digestHex string) error { + _, err := resolveTag(p, repository, tag) + if err == nil { + return nil + } + if !errors.Is(err, ErrNotFound) { + return err + } + return createTagSymlink(p, repository, tag, digestHex) +} + +func listTags(p *paths.Paths, repository string) ([]string, error) { + dirs := []string{filepath.Join(p.ImageRepositoriesDir(), repository), p.ImageRepositoryDir(repository)} + seen := make(map[string]struct{}) + tags := make([]string, 0) + for _, repoDir := range dirs { + entries, err := os.ReadDir(repoDir) + if err != nil { + if os.IsNotExist(err) { + continue + } + return nil, fmt.Errorf("read repository directory: %w", err) + } + for _, entry := range entries { + if entry.Type()&os.ModeSymlink == 0 { + continue + } + if _, ok := seen[entry.Name()]; ok { + continue + } + seen[entry.Name()] = struct{}{} + tags = append(tags, entry.Name()) + } + } + return tags, nil +} + +func promoteLegacyImages(p *paths.Paths) { + imagesDir := p.ImagesDir() + type legacyRef struct { + repository string + digestHex string + } + refs := make([]legacyRef, 0) + err := filepath.Walk(imagesDir, func(path string, info os.FileInfo, err error) error { + if err != nil || info.IsDir() || info.Name() != "metadata.json" { + return nil + } + rel, relErr := filepath.Rel(imagesDir, path) + if relErr != nil { + return nil + } + parts := strings.Split(rel, string(filepath.Separator)) + if len(parts) < 3 || parts[0] == "content" || parts[0] == "repositories" { + return nil + } + refs = append(refs, legacyRef{ + repository: filepath.Join(parts[:len(parts)-2]...), + digestHex: parts[len(parts)-2], + }) + return nil + }) + if err != nil && !os.IsNotExist(err) { + fmt.Fprintf(os.Stderr, "Warning: failed to scan legacy images for promotion: %v\n", err) + return + } + + for _, ref := range refs { + layout := resolveImageLayout(p, ref.repository, ref.digestHex) + meta, readErr := readMetadataAt(layout) + if readErr != nil || meta.Status != StatusReady { + continue + } + if promoteErr := promoteImageToContent(p, ref.repository, ref.digestHex, meta); promoteErr != nil { + fmt.Fprintf(os.Stderr, "Warning: failed to promote legacy image %s@%s: %v\n", ref.repository, ref.digestHex, promoteErr) + } + } +} + +func listAllMetadata(p *paths.Paths) ([]*imageMetadata, error) { + imagesDir := p.ImagesDir() + seen := make(map[string]struct{}) + contentDigests := make(map[string]struct{}) + taggedDigests := make(map[string]struct{}) + taggedContentDigests := make(map[string]struct{}) + metadataRefs := make([]metadataReference, 0) + seenMetadataRefs := make(map[string]struct{}) + metas := make([]*imageMetadata, 0) + + err := filepath.Walk(imagesDir, func(path string, info os.FileInfo, err error) error { + if err != nil { + return nil // Skip errors + } + rel, err := filepath.Rel(imagesDir, path) + if err != nil { + return nil + } + parts := strings.Split(rel, string(filepath.Separator)) + if len(parts) > 0 && parts[0] == "content" { + if info.IsDir() { + return nil + } + if info.Name() != "metadata.json" { + return nil + } + digestHex := filepath.Base(filepath.Dir(path)) + contentDigests[digestHex] = struct{}{} + return nil + } + + switch { + case info.Mode()&os.ModeSymlink != 0: + digestHex, err := os.Readlink(path) + if err != nil { + return nil // Skip invalid symlinks + } + digestHex = filepath.Base(digestHex) + + var repository, tag string + if len(parts) > 1 && parts[0] == "repositories" { + repository = filepath.Join(parts[1 : len(parts)-1]...) + tag = parts[len(parts)-1] + } else { + repository = filepath.Dir(rel) + tag = filepath.Base(path) + } + return appendMetadataForTag(p, repository, tag, digestHex, seen, taggedDigests, taggedContentDigests, &metas) + case !info.IsDir() && info.Name() == "metadata.json": + digestHex := filepath.Base(filepath.Dir(path)) + repository, err := filepath.Rel(imagesDir, filepath.Dir(filepath.Dir(path))) + if err != nil { + return nil + } + key := repository + "@" + digestHex + if _, ok := seenMetadataRefs[key]; !ok { + seenMetadataRefs[key] = struct{}{} + metadataRefs = append(metadataRefs, metadataReference{repository: repository, digestHex: digestHex}) + } + return nil + default: + return nil + } + }) + + if err != nil && !os.IsNotExist(err) { + return nil, fmt.Errorf("walk images directory: %w", err) + } + seenDigests := make(map[string]struct{}, len(metas)) + for _, ref := range metadataRefs { + if _, tagged := taggedDigests[ref.repository+"@"+ref.digestHex]; tagged { + continue + } + appendMetadataIfNew(p, ref.repository, ref.digestHex, seen, &metas) + seenDigests[ref.digestHex] = struct{}{} + } + for digestHex := range contentDigests { + if _, found := taggedContentDigests[digestHex]; found { + continue + } + if _, found := seenDigests[digestHex]; found { + continue + } + appendContentMetadataIfNew(p, digestHex, seen, &metas) + } + + return metas, nil +} + +type metadataReference struct { + repository string + digestHex string +} + +func appendMetadataIfNew(p *paths.Paths, repository, digestHex string, seen map[string]struct{}, metas *[]*imageMetadata) { + key := repository + "@" + digestHex + if _, ok := seen[key]; ok { + return + } + + meta, err := readMetadata(p, repository, digestHex) + if err != nil { + return + } + + seen[key] = struct{}{} + *metas = append(*metas, meta) +} + +func appendContentMetadataIfNew(p *paths.Paths, digestHex string, seen map[string]struct{}, metas *[]*imageMetadata) { + key := "@" + digestHex + if _, ok := seen[key]; ok { + return + } + meta, err := readContentMetadata(p, digestHex) + if err != nil { + return + } + seen[key] = struct{}{} + *metas = append(*metas, meta) +} + +func appendMetadataForTag(p *paths.Paths, repository, tag, digestHex string, seen, taggedDigests, taggedContentDigests map[string]struct{}, metas *[]*imageMetadata) error { + tagKey := repository + ":" + tag + if _, ok := seen[tagKey]; ok { + return nil + } + meta, err := readMetadata(p, repository, digestHex) + if err != nil { + return nil + } + meta.Name = repository + ":" + tag + seen[tagKey] = struct{}{} + taggedDigests[repository+"@"+digestHex] = struct{}{} + taggedContentDigests[digestHex] = struct{}{} + *metas = append(*metas, meta) + return nil +} + +func deleteTag(p *paths.Paths, repository, tag string) error { + pathsToRemove := []string{ + p.ImageRepositoryTagSymlink(repository, tag), + p.ImageTagSymlink(repository, tag), + } + found := false + for _, linkPath := range pathsToRemove { + if _, err := os.Lstat(linkPath); err != nil { + if os.IsNotExist(err) { + continue + } + return fmt.Errorf("stat symlink: %w", err) + } + found = true + if err := os.Remove(linkPath); err != nil { + return fmt.Errorf("remove symlink: %w", err) + } + } + if !found { + return ErrNotFound + } + return nil +} + +func tagsForDigest(p *paths.Paths, repository, digestHex string) ([]string, error) { + tags, err := listTags(p, repository) + if err != nil { + return nil, err + } + matched := make([]string, 0) + for _, tag := range tags { + target, err := resolveTag(p, repository, tag) + if err == nil && target == digestHex { + matched = append(matched, tag) + } + } + return matched, nil +} + +func countTagsForDigest(p *paths.Paths, repository, digestHex string) (int, error) { + tags, err := tagsForDigest(p, repository, digestHex) + return len(tags), err +} + +func deleteTagsForDigest(p *paths.Paths, repository, digestHex string) error { + tags, err := tagsForDigest(p, repository, digestHex) + if err != nil { + return err + } + for _, tag := range tags { + if err := deleteTag(p, repository, tag); err != nil && !errors.Is(err, ErrNotFound) { + return err + } + } + return nil +} + +func contentTagCount(p *paths.Paths, digestHex string) (int, error) { + root := p.ImageRepositoriesDir() + count := 0 + err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error { + if err != nil || info.IsDir() || info.Mode()&os.ModeSymlink == 0 { + return nil + } + rel, err := filepath.Rel(root, path) + if err != nil { + return nil + } + parts := strings.Split(rel, string(filepath.Separator)) + if len(parts) < 2 { + return nil + } + repository := filepath.Join(parts[:len(parts)-1]...) + if target, err := resolveTag(p, repository, parts[len(parts)-1]); err == nil && target == digestHex { + count++ + } + return nil + }) + if err != nil && !os.IsNotExist(err) { + return 0, fmt.Errorf("walk content tags: %w", err) + } + return count, nil +} + +func contentPullInProgress(p *paths.Paths, digestHex string) bool { + status, ok := metadataStatus(p.ImageContentMetadata(digestHex)) + return ok && (status == StatusPending || status == StatusPulling || status == StatusConverting) +} + +func contentIsDigestOnly(p *paths.Paths, digestHex string) bool { + meta, err := readContentMetadata(p, digestHex) + if err != nil { + return false + } + ref, err := ParseNormalizedRef(meta.Name) + return err == nil && ref.IsDigest() +} + +func removeDigestIfUnreferenced(p *paths.Paths, repository, digestHex string, preserveDigestOnly bool) error { + contentDir := p.ImageContentDir(digestHex) + contentExists := false + if _, err := os.Stat(contentDir); err == nil { + contentExists = true + } else if !os.IsNotExist(err) { + return fmt.Errorf("stat content digest directory: %w", err) + } + + if err := os.RemoveAll(p.ImageDigestDir(repository, digestHex)); err != nil { + return fmt.Errorf("remove legacy digest directory: %w", err) + } + if !contentExists { + return nil + } + + tagCount, err := contentTagCount(p, digestHex) + if err != nil { + return err + } + if tagCount > 0 || contentPullInProgress(p, digestHex) || (preserveDigestOnly && contentIsDigestOnly(p, digestHex)) { + return nil + } + + if err := os.RemoveAll(contentDir); err != nil { + return fmt.Errorf("remove content digest directory: %w", err) + } + return nil +} diff --git a/lib/images/storage_test.go b/lib/images/storage_test.go index f8ec44dcd..b34c47d54 100644 --- a/lib/images/storage_test.go +++ b/lib/images/storage_test.go @@ -1,41 +1,326 @@ package images import ( + "os" + "path/filepath" "testing" "time" + "github.com/kernel/hypeman/lib/paths" + "github.com/kernel/hypeman/lib/tags" "github.com/stretchr/testify/require" ) -func TestImageMetadataToImage_ClonesMetadata(t *testing.T) { - createdAt := time.Now().UTC().Truncate(time.Second) - source := &imageMetadata{ - Name: "docker.io/library/alpine:latest", - Digest: "sha256:abc", +func TestDeleteTagRemovesBothLayoutReferences(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + tag := "latest" + digest := "abababababababababababababababababababababababababababababababab" + legacyPath := p.ImageTagSymlink(repository, tag) + contentPath := p.ImageRepositoryTagSymlink(repository, tag) + require.NoError(t, os.MkdirAll(filepath.Dir(legacyPath), 0o755)) + require.NoError(t, os.MkdirAll(filepath.Dir(contentPath), 0o755)) + require.NoError(t, os.Symlink(digest, legacyPath)) + target, err := filepath.Rel(filepath.Dir(contentPath), p.ImageContentDir(digest)) + require.NoError(t, err) + require.NoError(t, os.Symlink(target, contentPath)) + + require.NoError(t, deleteTag(p, repository, tag)) + _, err = os.Lstat(legacyPath) + require.ErrorIs(t, err, os.ErrNotExist) + _, err = os.Lstat(contentPath) + require.ErrorIs(t, err, os.ErrNotExist) +} + +func TestLegacyImageIsNotShadowedByContentMetadata(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd" + legacyMeta := &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusReady, + } + legacyDir := p.ImageDigestDir(repository, digest) + require.NoError(t, os.MkdirAll(legacyDir, 0o755)) + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), legacyMeta)) + require.NoError(t, os.WriteFile(p.ImageDigestPath(repository, digest), []byte("legacy rootfs"), 0o644)) + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), &imageMetadata{ + Name: "docker.io/library/busybox:latest", + Digest: "sha256:" + digest, + Status: StatusFailed, + })) + + meta, err := readMetadata(p, repository, digest) + require.NoError(t, err) + require.Equal(t, StatusReady, meta.Status) + require.Equal(t, p.ImageDigestPath(repository, digest), digestPath(p, repository, digest)) +} + +func TestListAllMetadataDeduplicatesDualLayouts(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "9999999999999999999999999999999999999999999999999999999999999999" + meta := &imageMetadata{ + Name: repository + "@sha256:" + digest, + Digest: "sha256:" + digest, + Status: StatusReady, + } + + require.NoError(t, os.MkdirAll(p.ImageDigestDir(repository, digest), 0o755)) + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), meta)) + require.NoError(t, os.WriteFile(p.ImageDigestPath(repository, digest), []byte("legacy rootfs"), 0o644)) + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("content rootfs"), 0o644)) + + metas, err := listAllMetadata(p) + require.NoError(t, err) + require.Len(t, metas, 1) + require.Equal(t, meta.Digest, metas[0].Digest) +} + +func TestFailedLegacyImageUsesReadyContent(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee" + legacyDir := p.ImageDigestDir(repository, digest) + require.NoError(t, os.MkdirAll(legacyDir, 0o755)) + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusFailed, + })) + require.NoError(t, os.WriteFile(p.ImageDigestPath(repository, digest), []byte("legacy rootfs"), 0o644)) + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), &imageMetadata{ + Name: "registry.example.com/app:latest", + Digest: "sha256:" + digest, Status: StatusReady, - Tags: map[string]string{"team": "backend", "env": "staging"}, - SizeBytes: 123, - CreatedAt: createdAt, + SizeBytes: 13, + })) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("content rootfs"), 0o644)) + + meta, err := readMetadata(p, repository, digest) + require.NoError(t, err) + require.Equal(t, StatusReady, meta.Status) + require.Equal(t, p.ImageContentDir(digest), digestDir(p, repository, digest)) + require.Equal(t, p.ImageContentPath(digest), digestPath(p, repository, digest)) +} + +func TestReadyContentDoesNotFallBackToLegacyDisk(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff" + legacyDir := p.ImageDigestDir(repository, digest) + require.NoError(t, os.MkdirAll(legacyDir, 0o755)) + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusReady, + })) + require.NoError(t, os.WriteFile(p.ImageDigestPath(repository, digest), []byte("legacy rootfs"), 0o644)) + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusReady, + })) + + require.Equal(t, p.ImageContentMetadata(digest), metadataPath(p, repository, digest)) + require.Equal(t, p.ImageContentPath(digest), digestPath(p, repository, digest)) + _, err := readMetadata(p, repository, digest) + require.ErrorContains(t, err, "disk image missing") +} + +func TestWriteMetadataUsesContentWhenLegacyDirectoryIsEmpty(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + require.NoError(t, os.MkdirAll(p.ImageDigestDir(repository, digest), 0o755)) + require.NoError(t, writeMetadata(p, repository, digest, &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusPending, + })) + require.FileExists(t, p.ImageContentMetadata(digest)) + _, err := os.Stat(p.ImageMetadata(repository, digest)) + require.ErrorIs(t, err, os.ErrNotExist) +} + +func TestContentLayoutResolvesDiskByDigest(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "1212121212121212121212121212121212121212121212121212121212121212" + meta := &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusReady, + } + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("rootfs"), 0o644)) + + got, err := GetDiskPath(p, repository+":latest", "sha256:"+digest) + require.NoError(t, err) + require.Equal(t, p.ImageContentPath(digest), got) +} + +func TestSharedContentKeepsReferenceTags(t *testing.T) { + p := paths.New(t.TempDir()) + digest := "abababababababababababababababababababababababababababababababab" + first := "docker.io/library/alpine:latest" + second := "registry.example.com/app:v1" + meta := &imageMetadata{ + Name: first, + Digest: "sha256:" + digest, + Status: StatusReady, + References: map[string]tags.Tags{ + first: {"team": "one"}, + second: {"team": "two"}, + }, } + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("rootfs"), 0o644)) + require.NoError(t, createTagSymlink(p, "docker.io/library/alpine", "latest", digest)) + require.NoError(t, createTagSymlink(p, "registry.example.com/app", "v1", digest)) - img := source.toImage() - require.Equal(t, source.Name, img.Name) - require.Equal(t, source.Digest, img.Digest) - require.Equal(t, map[string]string{"team": "backend", "env": "staging"}, img.Tags) - require.NotNil(t, img.SizeBytes) - require.Equal(t, int64(123), *img.SizeBytes) + mgr := &manager{paths: p} + image, err := mgr.GetImage(nil, second) + require.NoError(t, err) + require.Equal(t, tags.Tags{"team": "two"}, image.Tags) - source.Tags["team"] = "mutated" - require.Equal(t, "backend", img.Tags["team"]) + images, err := mgr.ListImages(nil) + require.NoError(t, err) + got := make(map[string]tags.Tags, len(images)) + for _, image := range images { + got[image.Name] = image.Tags + } + require.Equal(t, map[string]tags.Tags{ + first: {"team": "one"}, + second: {"team": "two"}, + }, got) } -func TestImageMetadataToImage_EmptyMetadataOmitted(t *testing.T) { - img := (&imageMetadata{ +func TestPendingTagClaimSurvivesSharedBuild(t *testing.T) { + p := paths.New(t.TempDir()) + const repository = "docker.io/library/alpine" + const otherRepository = "registry.example.com/app" + const tag = "v1" + const digest = "cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd" + const previous = "efefefefefefefefefefefefefefefefefefefefefefefefefefefefefefefef" + + require.NoError(t, createTagSymlink(p, otherRepository, tag, previous)) + meta := &imageMetadata{ + Name: repository + ":latest", Digest: "sha256:" + digest, + Status: StatusPending, RequestedTag: "latest", + } + normalized, err := ParseNormalizedRef(otherRepository + ":" + tag) + require.NoError(t, err) + m := &manager{paths: p, tagGenerations: make(map[string]uint64)} + ref := NewResolvedRef(normalized, "sha256:"+digest) + require.NoError(t, m.recordPendingTag(meta, ref)) + require.Len(t, meta.TagClaims, 1) + claim := meta.TagClaims[0] + require.Equal(t, previous, claim.PreviousTagDigest) + require.Equal(t, previous, mustResolveTag(t, p, otherRepository, tag)) + + meta.Status = StatusReady + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("rootfs"), 0o644)) + m.claimTag(claim.Repository, digest, claim.Tag, claim.PreviousTagDigest, claim.TagGeneration, false) + require.Equal(t, digest, mustResolveTag(t, p, otherRepository, tag)) +} + +func mustResolveTag(t *testing.T, p *paths.Paths, repository, tag string) string { + t.Helper() + resolved, err := resolveTag(p, repository, tag) + require.NoError(t, err) + return resolved +} + +func TestListAllMetadataContentLayout(t *testing.T) { + p := paths.New(t.TempDir()) + digest := "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + meta := &imageMetadata{ Name: "docker.io/library/alpine:latest", - Digest: "sha256:abc", - Status: StatusPending, + Digest: "sha256:" + digest, + Status: StatusReady, + SizeBytes: 5, CreatedAt: time.Now().UTC(), - }).toImage() + } + require.NoError(t, writeMetadataFile(p.ImageContentMetadata(digest), meta)) + require.NoError(t, os.WriteFile(p.ImageContentPath(digest), []byte("rootfs"), 0o644)) + linkPath := p.ImageRepositoryTagSymlink("docker.io/library/alpine", "latest") + require.NoError(t, os.MkdirAll(filepath.Dir(linkPath), 0o755)) + target, err := filepath.Rel(filepath.Dir(linkPath), p.ImageContentDir(digest)) + require.NoError(t, err) + require.NoError(t, os.Symlink(target, linkPath)) + require.NoError(t, createTagSymlink(p, "docker.io/library/alpine", "stable", digest)) + _, err = os.Lstat(p.ImageRepositoryTagSymlink("docker.io/library/alpine", "stable")) + require.NoError(t, err) + + metas, err := listAllMetadata(p) + require.NoError(t, err) + require.Len(t, metas, 2) + names := []string{metas[0].Name, metas[1].Name} + require.ElementsMatch(t, []string{ + "docker.io/library/alpine:latest", + "docker.io/library/alpine:stable", + }, names) +} + +func TestPromoteLegacyImagesMovesContentAndTags(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + tag := "latest" + digest := "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc" + + legacyDir := p.ImageDigestDir(repository, digest) + require.NoError(t, os.MkdirAll(legacyDir, 0o755)) + meta := &imageMetadata{ + Name: repository + ":" + tag, + Digest: "sha256:" + digest, + Status: StatusReady, + SizeBytes: 7, + CreatedAt: time.Now().UTC(), + } + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), meta)) + require.NoError(t, os.WriteFile(p.ImageDigestPath(repository, digest), []byte("rootfs!"), 0o644)) + tagPath := p.ImageTagSymlink(repository, tag) + require.NoError(t, os.MkdirAll(filepath.Dir(tagPath), 0o755)) + require.NoError(t, os.Symlink(digest, tagPath)) + + promoteLegacyImages(p) + + contentMeta, err := readContentMetadata(p, digest) + require.NoError(t, err) + require.Equal(t, StatusReady, contentMeta.Status) + data, err := os.ReadFile(p.ImageContentPath(digest)) + require.NoError(t, err) + require.Equal(t, "rootfs!", string(data)) + + require.FileExists(t, p.ImageDigestPath(repository, digest)) + + resolved, err := resolveTag(p, repository, tag) + require.NoError(t, err) + require.Equal(t, digest, resolved) +} + +func TestPromoteLegacyImagesSkipsNonReady(t *testing.T) { + p := paths.New(t.TempDir()) + repository := "docker.io/library/alpine" + digest := "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee" + + legacyDir := p.ImageDigestDir(repository, digest) + require.NoError(t, os.MkdirAll(legacyDir, 0o755)) + require.NoError(t, writeMetadataFile(p.ImageMetadata(repository, digest), &imageMetadata{ + Name: repository + ":latest", + Digest: "sha256:" + digest, + Status: StatusConverting, + CreatedAt: time.Now().UTC(), + })) + + promoteLegacyImages(p) - require.Nil(t, img.Tags) + _, err := os.Stat(p.ImageContentMetadata(digest)) + require.True(t, os.IsNotExist(err), "non-ready legacy image must not be promoted") + _, err = os.Stat(legacyDir) + require.NoError(t, err, "non-ready legacy tree must be preserved") } diff --git a/lib/paths/paths.go b/lib/paths/paths.go index add23ff52..814dc1432 100644 --- a/lib/paths/paths.go +++ b/lib/paths/paths.go @@ -142,6 +142,35 @@ func (p *Paths) UFFDSessionsDir(versionKey string) string { // Image path methods +// ImageContentDir returns the directory for content-addressed image data. +func (p *Paths) ImageContentDir(digestHex string) string { + return filepath.Join(p.dataDir, "images", "content", digestHex) +} + +// ImageContentPath returns the path to a content-addressed rootfs disk file. +func (p *Paths) ImageContentPath(digestHex string) string { + ext := "erofs" + if runtime.GOOS == "darwin" { + ext = "ext4" + } + return filepath.Join(p.ImageContentDir(digestHex), "rootfs."+ext) +} + +// ImageContentMetadata returns the path to metadata for content-addressed image data. +func (p *Paths) ImageContentMetadata(digestHex string) string { + return filepath.Join(p.ImageContentDir(digestHex), "metadata.json") +} + +// ImageRepositoriesDir returns the root directory for repository tag references. +func (p *Paths) ImageRepositoriesDir() string { + return filepath.Join(p.dataDir, "images", "repositories") +} + +// ImageRepositoryTagSymlink returns the path to a tag reference in the new layout. +func (p *Paths) ImageRepositoryTagSymlink(repository, tag string) string { + return filepath.Join(p.ImageRepositoriesDir(), repository, tag) +} + // ImageDigestDir returns the directory for a specific image digest. func (p *Paths) ImageDigestDir(repository, digestHex string) string { return filepath.Join(p.dataDir, "images", repository, digestHex) diff --git a/lib/snapshot/testsupport/images.go b/lib/snapshot/testsupport/images.go index d033cd0ff..0608368c3 100644 --- a/lib/snapshot/testsupport/images.go +++ b/lib/snapshot/testsupport/images.go @@ -73,15 +73,19 @@ func EnsureImageReady(t *testing.T, ctx context.Context, p *paths.Paths, imageMa digestHex := strings.TrimPrefix(cached.Digest, "sha256:") require.NotEmpty(t, digestHex) - srcDigestDir := cachePaths.ImageDigestDir(ref.Repository(), digestHex) - dstDigestDir := p.ImageDigestDir(ref.Repository(), digestHex) + srcDiskPath, err := images.GetDiskPath(cachePaths, waitName, cached.Digest) + require.NoError(t, err) + srcDigestDir := filepath.Dir(srcDiskPath) + dstDigestDir := p.ImageContentDir(digestHex) require.NoError(t, copyDirWithHardlinks(srcDigestDir, dstDigestDir)) if ref.Tag() != "" { - linkPath := p.ImageTagSymlink(ref.Repository(), ref.Tag()) + linkPath := p.ImageRepositoryTagSymlink(ref.Repository(), ref.Tag()) + target, err := filepath.Rel(filepath.Dir(linkPath), dstDigestDir) + require.NoError(t, err) require.NoError(t, os.MkdirAll(filepath.Dir(linkPath), 0755)) _ = os.Remove(linkPath) - require.NoError(t, os.Symlink(digestHex, linkPath)) + require.NoError(t, os.Symlink(target, linkPath)) } reference := ref.Tag()