diff --git a/lib/images/disk_usage_test.go b/lib/images/disk_usage_test.go index 0de352d48..bd6056fef 100644 --- a/lib/images/disk_usage_test.go +++ b/lib/images/disk_usage_test.go @@ -9,6 +9,8 @@ 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)) @@ -21,6 +23,8 @@ func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForMalformedMetadata } func TestTotalReadyImageBytesFromMetadata_DeduplicatesHardLinkedAliases(t *testing.T) { + t.Parallel() + imagesDir := t.TempDir() sourceDir := filepath.Join(imagesDir, "source", "digest") targetDir := filepath.Join(imagesDir, "target", "digest") @@ -41,6 +45,8 @@ func TestTotalReadyImageBytesFromMetadata_DeduplicatesHardLinkedAliases(t *testi } func TestTotalReadyImageBytesFromMetadata_DeduplicatesMalformedAliases(t *testing.T) { + t.Parallel() + imagesDir := t.TempDir() malformedDir := filepath.Join(imagesDir, "a-malformed", "digest") validDir := filepath.Join(imagesDir, "b-valid", "digest") @@ -59,6 +65,8 @@ func TestTotalReadyImageBytesFromMetadata_DeduplicatesMalformedAliases(t *testin } func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForReadyImageWithoutSize(t *testing.T) { + t.Parallel() + 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 06e6cbbe0..27b017dbe 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -77,6 +77,7 @@ type manager struct { createMu sync.Mutex diskUsageMu sync.RWMutex tagGenerations map[string]uint64 + requestedTags map[string]string // newest pull's digest per requested tag diskUsageLoaded bool readyImageBytes int64 ociCacheBytes int64 @@ -105,6 +106,7 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, readySubscribers: make(map[string][]chan StatusEvent), tagGenerations: make(map[string]uint64), + requestedTags: make(map[string]string), } // Initialize metrics if meter is provided @@ -117,7 +119,11 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma } m.RecoverInterruptedBuilds() - promoteLegacyImages(m.paths) + // Promotion hardlinks disks and repoints tags, which can take a while on a + // large legacy store; run it in the background so the manager is usable + // immediately. Failures only warn so a partial migration never blocks + // startup. + go promoteLegacyImages(p) return m, nil } @@ -212,7 +218,7 @@ func (m *manager) CreateImage(ctx context.Context, req CreateImageRequest) (*Ima m.createMu.Lock() defer m.createMu.Unlock() - if img, found, err := m.reuseExistingImage(ref, req.Credentials, req.Tags); found || err != nil { + if img, found, err := m.reuseExistingImage(ref, req.Credentials); found || err != nil { return img, err } return m.createAndQueueImage(ref, req, platform) @@ -244,14 +250,31 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest m.createMu.Lock() defer m.createMu.Unlock() - if img, found, err := m.reuseExistingImage(ref, nil, nil); found || err != nil { - return img, err + // 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 by falling through to + // re-queue the build. + if meta.Status == StatusFailed { + if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { + return nil, fmt.Errorf("remove failed image: %w", err) + } + } else { + if err := m.claimTagForStatus(meta, ref); err != nil { + return nil, fmt.Errorf("create image tag: %w", err) + } + img := meta.toImageFor(ref.String()) + if meta.Status == StatusPending { + img.QueuePosition = m.queue.GetPosition(meta.Digest) + } + return img, nil + } } + // 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) { +func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthConfig) (*Image, bool, error) { meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()) if err != nil { return nil, false, nil @@ -262,95 +285,25 @@ func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthCo } return nil, false, nil } - if err := m.updateExistingReference(meta, ref, resourceTags); err != nil { - return nil, true, fmt.Errorf("update image reference: %w", err) + if ref.Tag() != "" { + if meta.Status != StatusReady { + // A pending pull with different credentials does not get to point + // the tag at its digest. + if !m.inflightCredentialsMatch(ref.Digest(), credentials) { + return nil, true, fmt.Errorf("%w: retry after the current pull completes", ErrCredentialConflict) + } + } + if err := m.claimTagForStatus(meta, ref); err != nil { + return nil, true, fmt.Errorf("create image tag: %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 @@ -442,6 +395,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, tagGeneration := uint64(0) if ref.Tag() != "" { tagGeneration = m.nextTagGeneration(ref.Repository(), ref.Tag()) + m.trackRequestedTag(ref.Repository(), ref.Tag(), ref.DigestHex()) } meta := &imageMetadata{ Name: ref.String(), @@ -455,16 +409,19 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, 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 { + if ref.Tag() != "" { + m.revertTagGeneration(ref.Repository(), ref.Tag()) + } 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 { + m.revertTagGeneration(ref.Repository(), ref.Tag()) return nil, fmt.Errorf("create pending image tag: %w", err) } } @@ -542,8 +499,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials // Another build completed first; last-pull-wins repoints the tag. if ref.Tag() != "" { m.createMu.Lock() - m.nextTagGeneration(ref.Repository(), ref.Tag()) - err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + err := m.claimReadyTag(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) @@ -557,6 +513,8 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials m.updateStatusByDigest(ref, StatusConverting, nil, buildID) diskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk + // Keep the temporary filesystem beside its final path so finalization stays + // atomic even when system/builds and images are on different filesystems. diskTempPath := diskPath + ".tmp-" + buildID defer os.Remove(diskTempPath) // Use default image format (erofs on Linux, ext4 on Darwin) @@ -583,21 +541,18 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } 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() + layout := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()) + // Read current metadata to preserve request info and reject stale builds. - meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()) + meta, err := readMetadataAt(layout) if err != nil || meta.BuildID != buildID { return errStaleBuild } - finalDiskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk - if err := installAtomically(finalDiskPath, func(path string) error { + if err := installAtomically(layout.disk, func(path string) error { return os.Rename(diskTempPath, path) }); err != nil { return fmt.Errorf("install image disk: %w", err) @@ -613,6 +568,17 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i return err } + // Persist the manifest content model beside the shared content so later + // stages can recompose the image from per-layer artifacts and GC can tell + // which OCI blobs are still referenced. + if result.Manifest != nil { + model := *result.Manifest + model.Platform = actualPlatform.String() + if err := writeManifestModel(m.paths, ref.DigestHex(), &model); err != nil { + return fmt.Errorf("write manifest model: %w", err) + } + } + meta.Status = StatusReady meta.Error = nil meta.Platform = actualPlatform.String() @@ -623,44 +589,16 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i meta.Labels = result.Metadata.Labels meta.WorkingDir = result.Metadata.WorkingDir - if err := writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta); err != nil { + if err := writeMetadataFile(layout.metadata, meta); err != nil { return fmt.Errorf("write final metadata: %w", err) } m.notifyReady(ref.DigestHex(), StatusReady, nil) - m.claimImageTags(ref, meta) + m.claimRequestedTag(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" @@ -705,7 +643,8 @@ func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err erro m.createMu.Lock() defer m.createMu.Unlock() - meta, readErr := readMetadata(m.paths, ref.Repository(), ref.DigestHex()) + layout := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()) + meta, readErr := readMetadataAt(layout) if readErr != nil || meta.BuildID != buildID { return } @@ -716,13 +655,18 @@ func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err erro meta.Error = &errorMsg } - writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta) + writeMetadataFile(layout.metadata, meta) // Notify while holding createMu so a delete/recreate cannot race the // metadata write and receive a terminal event for the old build. if status == StatusReady || status == StatusFailed { m.notifyReady(ref.DigestHex(), status, err) } + // A failed pull releases its tag claim so an older in-flight pull of the + // same tag can still repoint it. + if status == StatusFailed && meta.RequestedTag != "" { + m.releaseTagGeneration(ref.Repository(), meta.RequestedTag, meta.TagGeneration) + } } func (m *manager) RecoverInterruptedBuilds() { @@ -730,12 +674,12 @@ 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) }) + m.restoreTagState(metas) seenDigests := make(map[string]struct{}) for _, meta := range metas { @@ -808,6 +752,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := deleteTagsForDigest(m.paths, repository, digestHex); err != nil { return err } + m.pruneTagGenerations() if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, false); err != nil { return err } @@ -816,7 +761,6 @@ 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) @@ -828,6 +772,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := deleteTag(m.paths, repository, tag); err != nil { return err } + m.pruneTagGenerations() // Check if the digest is now orphaned (no other tags reference it) count, err := countTagsForDigest(m.paths, repository, digestHex) @@ -863,30 +808,6 @@ func (m *manager) TotalOCICacheBytes(ctx context.Context) (int64, error) { return ociCacheBytes, nil } -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 requestedTag != ref.Tag() || !strings.HasPrefix(meta.Name, ref.Repository()+":") { - continue - } - if newest == nil || newest.CreatedAt.Before(meta.CreatedAt) { - newest = meta - } - } - if newest == nil { - return nil - } - return newest.toImageFor(ref.String()) -} - // 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 { @@ -898,11 +819,8 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { if err != nil { return err } - if img.Status == StatusReady { - return nil - } - if img.Status == StatusFailed { - return conversionFailedErr(img.Error, nil) + if terminal, err := terminalImageError(img); terminal { + return err } digestHex := strings.TrimPrefix(img.Digest, "sha256:") @@ -915,11 +833,8 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { // Re-check after subscribing to close the race window img, err = m.GetImage(ctx, ref.Repository()+"@"+img.Digest) if err == nil { - if img.Status == StatusReady { - return nil - } - if img.Status == StatusFailed { - return conversionFailedErr(img.Error, nil) + if terminal, terminalErr := terminalImageError(img); terminal { + return terminalErr } } @@ -943,7 +858,7 @@ func (m *manager) waitForImage(ctx context.Context, name string, ref *Normalized for { var img *Image if !ref.IsDigest() { - img = m.findRequestedTagImage(ref) + img = m.requestedTagImage(ref) } if img == nil { img, lastErr = m.GetImage(ctx, name) @@ -962,6 +877,17 @@ func (m *manager) waitForImage(ctx context.Context, name string, ref *Normalized } } +func terminalImageError(img *Image) (bool, error) { + switch img.Status { + case StatusReady: + return true, nil + case StatusFailed: + return true, conversionFailedErr(img.Error, nil) + default: + return false, nil + } +} + 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 8abd90a65..13a6835e2 100644 --- a/lib/images/manager_test.go +++ b/lib/images/manager_test.go @@ -734,9 +734,9 @@ func TestDeleteAndRecreateDuringBuildTail(t *testing.T) { require.NoError(t, err) staleRef := NewResolvedRef(normalized, digestStr) m.updateStatusByDigest(staleRef, StatusFailed, errors.New("stale build"), firstMeta.BuildID) - staleResult, _, _, err := m.ociClient.extractOCIImageDetails(digestHex) + staleBundle, err := m.ociClient.extractOCIImageBundle(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: staleBundle.Meta}, 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/manifest_model.go b/lib/images/manifest_model.go new file mode 100644 index 000000000..b3e7dc25b --- /dev/null +++ b/lib/images/manifest_model.go @@ -0,0 +1,166 @@ +package images + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + + "github.com/kernel/hypeman/lib/paths" + digest "github.com/opencontainers/go-digest" + v1 "github.com/opencontainers/image-spec/specs-go/v1" +) + +// imageManifestModel is the persisted content model for one OCI manifest. It +// records everything needed to recompose the image from shared layer +// artifacts and to decide which OCI blobs are still referenced: the immutable +// manifest digest, the resolved platform, the image config with its diff ids, +// and the ordered layer descriptors. +type imageManifestModel struct { + SchemaVersion int `json:"schema_version"` + Digest string `json:"digest"` // manifest digest, sha256:... + MediaType string `json:"media_type,omitempty"` + Platform string `json:"platform"` // os/arch[/variant] + Config manifestConfigRef `json:"config"` + RootFSType string `json:"rootfs_type,omitempty"` + Layers []layerDescriptor `json:"layers"` // manifest order, base layer first +} + +const manifestModelSchemaVersion = 1 + +type manifestConfigRef struct { + Digest string `json:"digest"` // config blob digest, sha256:... + MediaType string `json:"media_type,omitempty"` + DiffIDs []string `json:"diff_ids,omitempty"` // uncompressed layer ids, matching Layers order +} + +// layerDescriptor describes one compressed layer blob exactly as it appears in +// the manifest, so the blob can be located in the shared OCI cache by digest. +type layerDescriptor struct { + Digest string `json:"digest"` // compressed blob digest, sha256:... + Size int64 `json:"size"` // compressed bytes + MediaType string `json:"media_type,omitempty"` + DiffID string `json:"diff_id,omitempty"` // uncompressed diff id from the image config +} + +// digestFromHex returns the full sha256 digest string for a bare hex value. +func digestFromHex(hex string) string { + if strings.HasPrefix(hex, "sha256:") { + return hex + } + return "sha256:" + hex +} + +// blobReferences returns every OCI blob digest the manifest depends on: the +// config blob plus all layer blobs. GC must keep these while the manifest is +// referenced. +func (m *imageManifestModel) blobReferences() []string { + refs := make([]string, 0, len(m.Layers)+1) + if m.Config.Digest != "" { + refs = append(refs, m.Config.Digest) + } + for _, layer := range m.Layers { + refs = append(refs, layer.Digest) + } + return refs +} + +func validateManifestModel(digestHex string, model *imageManifestModel) error { + if model == nil { + return fmt.Errorf("manifest model is nil") + } + if model.SchemaVersion != manifestModelSchemaVersion { + return fmt.Errorf("unsupported manifest model schema version: %d", model.SchemaVersion) + } + if model.Digest != digestFromHex(digestHex) { + return fmt.Errorf("manifest model digest %q does not match %q", model.Digest, digestFromHex(digestHex)) + } + if model.RootFSType != "" && model.RootFSType != "layers" { + return fmt.Errorf("unsupported manifest rootfs type: %q", model.RootFSType) + } + if err := validateManifestConfig(model); err != nil { + return err + } + return validateManifestLayers(model) +} + +func validateManifestConfig(model *imageManifestModel) error { + if model.Config.Digest == "" { + return fmt.Errorf("manifest model config digest is empty") + } + if _, err := parseSHA256Digest(model.Config.Digest); err != nil { + return fmt.Errorf("invalid manifest model config digest: %q", model.Config.Digest) + } + if model.Config.MediaType != "" && convertToOCIMediaType(model.Config.MediaType) != v1.MediaTypeImageConfig { + return fmt.Errorf("invalid manifest model config media type: %q", model.Config.MediaType) + } + if len(model.Config.DiffIDs) != len(model.Layers) { + return fmt.Errorf("manifest model has %d diff ids for %d layers", len(model.Config.DiffIDs), len(model.Layers)) + } + return nil +} + +func validateManifestLayers(model *imageManifestModel) error { + for i, layer := range model.Layers { + if _, err := parseSHA256Digest(layer.Digest); err != nil { + return fmt.Errorf("invalid manifest model layer %d digest: %q", i, layer.Digest) + } + diffID, err := parseSHA256Digest(model.Config.DiffIDs[i]) + if err != nil || layer.DiffID != diffID.String() { + return fmt.Errorf("invalid manifest model layer %d diff id", i) + } + if layer.Size < 0 { + return fmt.Errorf("invalid manifest model layer %d size: %d", i, layer.Size) + } + } + return nil +} + +func parseSHA256Digest(value string) (digest.Digest, error) { + parsed, err := digest.Parse(value) + if err != nil || parsed.Algorithm() != digest.SHA256 { + return "", fmt.Errorf("not a sha256 digest") + } + return parsed, nil +} + +func writeManifestModel(p *paths.Paths, digestHex string, model *imageManifestModel) error { + data, err := json.MarshalIndent(model, "", " ") + if err != nil { + return fmt.Errorf("marshal manifest model: %w", err) + } + return writeJSONAtomic(p.ImageContentManifestModel(digestHex), data) +} + +// readManifestModel loads the manifest model for a digest, if present. +// Missing models return (nil, nil): images converted before the manifest model +// existed only have a flattened rootfs. +func readManifestModel(p *paths.Paths, digestHex string) (*imageManifestModel, error) { + data, err := os.ReadFile(p.ImageContentManifestModel(digestHex)) + if err != nil { + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("read manifest model: %w", err) + } + var model imageManifestModel + if err := json.Unmarshal(data, &model); err != nil { + return nil, fmt.Errorf("unmarshal manifest model: %w", err) + } + if err := validateManifestModel(digestHex, &model); err != nil { + return nil, err + } + return &model, nil +} + +// writeJSONAtomic writes data to path via a temp file in the same directory +// followed by a rename, so readers never observe a partial document. +func writeJSONAtomic(path string, data []byte) error { + if err := installAtomically(path, func(tempPath string) error { + return os.WriteFile(tempPath, data, 0o644) + }); err != nil { + return fmt.Errorf("write %s: %w", filepath.Base(path), err) + } + return nil +} diff --git a/lib/images/manifest_model_test.go b/lib/images/manifest_model_test.go new file mode 100644 index 000000000..1137e3912 --- /dev/null +++ b/lib/images/manifest_model_test.go @@ -0,0 +1,242 @@ +package images + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "context" + "io" + "os" + "strings" + "testing" + "time" + + gcr "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/empty" + "github.com/google/go-containerregistry/pkg/v1/layout" + "github.com/google/go-containerregistry/pkg/v1/mutate" + "github.com/google/go-containerregistry/pkg/v1/tarball" + "github.com/google/go-containerregistry/pkg/v1/types" + "github.com/kernel/hypeman/lib/paths" + "github.com/stretchr/testify/require" +) + +// syntheticLayer builds a gzipped tar layer containing one file. +func syntheticLayer(t *testing.T, name, content string) gcr.Layer { + t.Helper() + + var buf bytes.Buffer + gzw := gzip.NewWriter(&buf) + tw := tar.NewWriter(gzw) + require.NoError(t, tw.WriteHeader(&tar.Header{ + Name: name, + Size: int64(len(content)), + Typeflag: tar.TypeReg, + Mode: 0644, + })) + _, err := tw.Write([]byte(content)) + require.NoError(t, err) + require.NoError(t, tw.Close()) + require.NoError(t, gzw.Close()) + + data := buf.Bytes() + layer, err := tarball.LayerFromOpener(func() (io.ReadCloser, error) { + return io.NopCloser(bytes.NewReader(data)), nil + }) + require.NoError(t, err) + return layer +} + +// writeSyntheticLayout writes img into a fresh OCI layout cache tagged with the +// image's own digest, mirroring pullToOCILayout. +func writeSyntheticLayout(t *testing.T, img gcr.Image) (*ociClient, string) { + t.Helper() + + client, err := newOCIClient(t.TempDir()) + require.NoError(t, err) + + digest, err := img.Digest() + require.NoError(t, err) + layoutPath, err := layout.Write(client.cacheDir, empty.Index) + require.NoError(t, err) + require.NoError(t, layoutPath.AppendImage(img, layout.WithAnnotations(map[string]string{ + "org.opencontainers.image.ref.name": digestToLayoutTag(digest.String()), + }))) + return client, digestToLayoutTag(digest.String()) +} + +func TestExtractManifestModel(t *testing.T) { + base := syntheticLayer(t, "base.txt", "base layer content") + top := syntheticLayer(t, "top.txt", "top layer content") + img, err := mutate.AppendLayers(empty.Image, base, top) + require.NoError(t, err) + img = mutate.MediaType(img, types.OCIManifestSchema1) + + client, layoutTag := writeSyntheticLayout(t, img) + + bundle, err := client.extractOCIImageBundle(layoutTag) + require.NoError(t, err) + model := bundle.Model + + manifest, err := img.Manifest() + require.NoError(t, err) + configFile, err := img.ConfigFile() + require.NoError(t, err) + configDigest, err := img.ConfigName() + require.NoError(t, err) + + require.Equal(t, "sha256:"+layoutTag, model.Digest) + require.Equal(t, manifest.MediaType, types.MediaType(model.MediaType)) + require.Equal(t, configDigest.String(), model.Config.Digest) + require.Len(t, model.Layers, 2) + require.Len(t, model.Config.DiffIDs, 2) + + // Layer order and pairing must match the manifest and config diff ids. + for i, desc := range model.Layers { + require.Equal(t, manifest.Layers[i].Digest.String(), desc.Digest) + require.Equal(t, manifest.Layers[i].Size, desc.Size) + require.Equal(t, configFile.RootFS.DiffIDs[i].String(), desc.DiffID) + require.Equal(t, model.Config.DiffIDs[i], desc.DiffID) + } + + // Digests are content addresses, so identical layer content dedupes. + require.NotEqual(t, model.Layers[0].Digest, model.Layers[1].Digest) +} + +func TestExtractManifestModelPlatform(t *testing.T) { + layer := syntheticLayer(t, "file.txt", "content") + img, err := mutate.AppendLayers(empty.Image, layer) + require.NoError(t, err) + cfgFile, err := img.ConfigFile() + require.NoError(t, err) + cfgFile = cfgFile.DeepCopy() + cfgFile.OS = "linux" + cfgFile.Architecture = "amd64" + img, err = mutate.ConfigFile(img, cfgFile) + require.NoError(t, err) + + client, layoutTag := writeSyntheticLayout(t, img) + + bundle, err := client.extractOCIImageBundle(layoutTag) + require.NoError(t, err) + model := bundle.Model + require.Equal(t, "linux/amd64", model.Platform) +} + +func TestManifestModelWriteReadRoundtrip(t *testing.T) { + p := paths.New(t.TempDir()) + digestHex := "ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01ab01" + + model := &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: "sha256:" + digestHex, + Platform: "linux/amd64", + Config: manifestConfigRef{ + Digest: "sha256:" + strings.Repeat("c", 64), + DiffIDs: []string{"sha256:" + strings.Repeat("d", 64), "sha256:" + strings.Repeat("e", 64)}, + }, + Layers: []layerDescriptor{ + {Digest: "sha256:" + strings.Repeat("f", 64), Size: 10, DiffID: "sha256:" + strings.Repeat("d", 64)}, + {Digest: "sha256:" + strings.Repeat("0", 64), Size: 20, DiffID: "sha256:" + strings.Repeat("e", 64)}, + }, + } + require.NoError(t, writeManifestModel(p, digestHex, model)) + + read, err := readManifestModel(p, digestHex) + require.NoError(t, err) + require.Equal(t, model, read) + + // No leftover temp files from the atomic write. + entries, err := os.ReadDir(p.ImageContentDir(digestHex)) + require.NoError(t, err) + require.Len(t, entries, 1) + require.Equal(t, "manifest.json", entries[0].Name()) +} + +func TestReadManifestModelRejectsInvalidSchema(t *testing.T) { + p := paths.New(t.TempDir()) + digestHex := strings.Repeat("a", 64) + require.NoError(t, os.MkdirAll(p.ImageContentDir(digestHex), 0755)) + require.NoError(t, os.WriteFile(p.ImageContentManifestModel(digestHex), []byte(`{"schema_version":1,"digest":"sha256:`+digestHex+`"}`), 0644)) + _, err := readManifestModel(p, digestHex) + require.ErrorContains(t, err, "config digest is empty") +} + +func TestReadManifestModelMissing(t *testing.T) { + p := paths.New(t.TempDir()) + model, err := readManifestModel(p, "deadbeef") + require.NoError(t, err) + require.Nil(t, model, "missing model must read as nil for pre-model images") +} + +func TestManifestModelBlobReferences(t *testing.T) { + model := &imageManifestModel{ + Config: manifestConfigRef{Digest: "sha256:cfg"}, + Layers: []layerDescriptor{ + {Digest: "sha256:layer0"}, + {Digest: "sha256:layer1"}, + }, + } + require.Equal(t, []string{"sha256:cfg", "sha256:layer0", "sha256:layer1"}, model.blobReferences()) +} + +func TestImportLocalImagePersistsManifestModel(t *testing.T) { + dataDir := t.TempDir() + p := paths.New(dataDir) + mgr, err := NewManager(p, 1, nil) + require.NoError(t, err) + m := mgr.(*manager) + + ctx := context.Background() + const repo = "kernel.local/test/manifest-model" + const tag = "v1" + + testImg := createTestDockerImage(t) + imgDigest, err := testImg.Digest() + require.NoError(t, err) + digestStr := imgDigest.String() + + cacheDir := p.SystemOCICache() + layoutPath, err := layout.Write(cacheDir, empty.Index) + require.NoError(t, err) + require.NoError(t, layoutPath.AppendImage(testImg, layout.WithAnnotations(map[string]string{ + "org.opencontainers.image.ref.name": digestToLayoutTag(digestStr), + }))) + + digestHex := digestToLayoutTag(digestStr) + events := make(chan StatusEvent, 2) + m.subscribeToReady(digestHex, events) + defer m.unsubscribeFromReady(digestHex, events) + + _, err = m.ImportLocalImage(ctx, repo, tag, digestStr) + require.NoError(t, err) + select { + case event := <-events: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatal("build did not become ready") + } + + model, err := readManifestModel(p, digestHex) + require.NoError(t, err) + require.NotNil(t, model, "ready image must persist a manifest model") + require.Equal(t, "sha256:"+digestHex, model.Digest) + + configFile, err := testImg.ConfigFile() + require.NoError(t, err) + manifest, err := testImg.Manifest() + require.NoError(t, err) + require.Len(t, model.Layers, len(manifest.Layers)) + for i, desc := range model.Layers { + require.Equal(t, manifest.Layers[i].Digest.String(), desc.Digest) + require.Equal(t, configFile.RootFS.DiffIDs[i].String(), desc.DiffID) + } + require.Equal(t, StatusReady, mustReadContentStatus(t, p, digestHex)) +} + +func mustReadContentStatus(t *testing.T, p *paths.Paths, digestHex string) string { + t.Helper() + meta, err := readContentMetadata(p, digestHex) + require.NoError(t, err) + return meta.Status +} diff --git a/lib/images/oci.go b/lib/images/oci.go index 33206d134..2a1a23223 100644 --- a/lib/images/oci.go +++ b/lib/images/oci.go @@ -194,6 +194,7 @@ func (c *ociClient) inspectDigestPlatformAuth(ctx context.Context, imageRef stri // pullResult contains the metadata and digest from pulling an image type pullResult struct { Metadata *containerMetadata + Manifest *imageManifestModel Digest string // sha256:abc123... CacheHit bool LayerCount int @@ -207,6 +208,14 @@ type imageBuildPhaseMeasurement struct { Status string } +// ociImageBundle is the extracted content of one OCI image in the cache. +type ociImageBundle struct { + Meta *containerMetadata + Model *imageManifestModel + LayerCount int + CompressedBytes int64 +} + func (r *pullResult) measure(phase string, operation func() error) error { start := time.Now() err := operation() @@ -260,20 +269,19 @@ func (c *ociClient) pullAndExportWithPlatformAuth(ctx context.Context, imageRef, // If cached, we skip the pull entirely // Extract metadata (from cache or freshly pulled) - var meta *containerMetadata - var layerCount int - var compressedBytes int64 + var bundle *ociImageBundle err := result.measure("metadata_extract", func() error { var err error - meta, layerCount, compressedBytes, err = c.extractOCIImageDetails(layoutTag) + bundle, err = c.extractOCIImageBundle(layoutTag) return err }) if err != nil { return result, fmt.Errorf("extract metadata: %w", err) } - result.Metadata = meta - result.LayerCount = layerCount - result.CompressedBytes = compressedBytes + result.Metadata = bundle.Meta + result.Manifest = bundle.Model + result.LayerCount = bundle.LayerCount + result.CompressedBytes = bundle.CompressedBytes // Unpack layers to the export directory if err := result.measure("layer_unpack", func() error { @@ -395,44 +403,31 @@ func imageByAnnotation(path layout.Path, layoutTag string) (gcr.Image, error) { return nil, fmt.Errorf("no image found with tag %s", layoutTag) } -// extractOCIMetadata reads metadata from OCI layout config.json -// Uses go-containerregistry which handles both Docker v2 and OCI v1 manifests. -func (c *ociClient) extractOCIMetadata(layoutTag string) (*containerMetadata, error) { - meta, _, _, err := c.extractOCIImageDetails(layoutTag) - return meta, err -} - -func (c *ociClient) extractOCIImageDetails(layoutTag string) (*containerMetadata, int, int64, error) { - // Open OCI layout using go-containerregistry (handles Docker v2 and OCI v1) +// extractOCIImageBundle reads metadata, the manifest content model, and layer +// stats from the cached OCI layout. Uses go-containerregistry which handles +// both Docker v2 and OCI v1 manifests. +func (c *ociClient) extractOCIImageBundle(layoutTag string) (*ociImageBundle, error) { path, err := layout.FromPath(c.cacheDir) if err != nil { - return nil, 0, 0, fmt.Errorf("open oci layout: %w", err) + return nil, fmt.Errorf("open oci layout: %w", err) } - - // Get the image by annotation tag from the layout img, err := imageByAnnotation(path, layoutTag) if err != nil { - return nil, 0, 0, fmt.Errorf("find image by tag %s: %w", layoutTag, err) + return nil, fmt.Errorf("find image by tag %s: %w", layoutTag, err) } - - // Get config file (go-containerregistry handles manifest format automatically) configFile, err := img.ConfigFile() if err != nil { - return nil, 0, 0, fmt.Errorf("get config file: %w", err) + return nil, fmt.Errorf("get config file: %w", err) } - manifest, err := img.Manifest() if err != nil { - return nil, 0, 0, fmt.Errorf("get manifest: %w", err) + return nil, fmt.Errorf("get manifest: %w", err) } - var compressedBytes int64 - for _, layer := range manifest.Layers { - compressedBytes += layer.Size + configDigest, err := img.ConfigName() + if err != nil { + return nil, fmt.Errorf("get config digest: %w", err) } - // Extract metadata from config. OS/Architecture/Variant come straight from - // the pulled image config, so they reflect the manifest actually fetched - // rather than what the caller requested. meta := &containerMetadata{ OS: configFile.OS, Architecture: configFile.Architecture, @@ -443,7 +438,6 @@ func (c *ociClient) extractOCIImageDetails(layoutTag string) (*containerMetadata Labels: make(map[string]string), WorkingDir: configFile.Config.WorkingDir, } - // Parse environment variables for _, env := range configFile.Config.Env { for i := 0; i < len(env); i++ { @@ -455,12 +449,49 @@ func (c *ociClient) extractOCIImageDetails(layoutTag string) (*containerMetadata } } } - for key, value := range configFile.Config.Labels { meta.Labels[key] = value } - return meta, len(manifest.Layers), compressedBytes, nil + model := &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: digestFromHex(layoutTag), + MediaType: string(manifest.MediaType), + Platform: Platform{ + OS: configFile.OS, + Architecture: configFile.Architecture, + Variant: configFile.Variant, + }.String(), + Config: manifestConfigRef{ + Digest: configDigest.String(), + MediaType: string(manifest.Config.MediaType), + DiffIDs: make([]string, 0, len(configFile.RootFS.DiffIDs)), + }, + RootFSType: configFile.RootFS.Type, + Layers: make([]layerDescriptor, 0, len(manifest.Layers)), + } + for _, diffID := range configFile.RootFS.DiffIDs { + model.Config.DiffIDs = append(model.Config.DiffIDs, diffID.String()) + } + var compressedBytes int64 + for i, descriptor := range manifest.Layers { + compressedBytes += descriptor.Size + layer := layerDescriptor{ + Digest: descriptor.Digest.String(), + Size: descriptor.Size, + MediaType: string(descriptor.MediaType), + } + if i < len(model.Config.DiffIDs) { + layer.DiffID = model.Config.DiffIDs[i] + } + model.Layers = append(model.Layers, layer) + } + return &ociImageBundle{ + Meta: meta, + Model: model, + LayerCount: len(manifest.Layers), + CompressedBytes: compressedBytes, + }, nil } // unpackLayers unpacks all OCI layers to a target directory using umoci diff --git a/lib/images/oci_test.go b/lib/images/oci_test.go index e076156ff..2cb765cc5 100644 --- a/lib/images/oci_test.go +++ b/lib/images/oci_test.go @@ -89,11 +89,11 @@ func TestExtractMetadataSucceedsOnBuildKitCache(t *testing.T) { // This succeeds because go-containerregistry doesn't validate config mediatype // The failure only happens in unpackLayers when umoci validates the config - meta, err := client.extractOCIMetadata("test-cache") - require.NoError(t, err, "extractOCIMetadata succeeds - go-containerregistry is lenient") + bundle, err := client.extractOCIImageBundle("test-cache") + require.NoError(t, err, "extractOCIImageBundle succeeds - go-containerregistry is lenient") // But the metadata will be empty/invalid since it's not a real OCI config - t.Logf("Got metadata (likely empty): %+v", meta) + t.Logf("Got metadata (likely empty): %+v", bundle.Meta) } // createBuildKitCacheLayout creates an OCI layout that mimics what BuildKit @@ -337,9 +337,10 @@ func TestDockerSaveTarballToOCILayoutRoundtrip(t *testing.T) { require.NoError(t, err) assert.True(t, client.existsInLayout(layoutTag), "image should exist in layout after AppendImage") - // Step 6: Verify extractOCIMetadata reads correct config - meta, err := client.extractOCIMetadata(layoutTag) + // Step 6: Verify extractOCIImageBundle reads correct config + bundle, err := client.extractOCIImageBundle(layoutTag) require.NoError(t, err) + meta := bundle.Meta assert.Equal(t, []string{"/usr/local/bin/guest-agent"}, meta.Entrypoint) assert.Equal(t, "/app", meta.WorkingDir) assert.Contains(t, meta.Env, "PATH") diff --git a/lib/images/storage.go b/lib/images/storage.go index 65d433c47..91ce646f7 100644 --- a/lib/images/storage.go +++ b/lib/images/storage.go @@ -14,53 +14,29 @@ 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"` - 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) + 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"` } 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 } @@ -216,24 +192,11 @@ func writeMetadata(p *paths.Paths, repository, digestHex string, meta *imageMeta } 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, "", " ") if err != nil { return fmt.Errorf("marshal metadata: %w", err) } - - tempPath := path + ".tmp" - if err := os.WriteFile(tempPath, data, 0644); err != nil { - return fmt.Errorf("write temp metadata: %w", err) - } - if err := os.Rename(tempPath, path); err != nil { - _ = os.Remove(tempPath) - return fmt.Errorf("rename metadata: %w", err) - } - return nil + return writeJSONAtomic(path, data) } func readMetadata(p *paths.Paths, repository, digestHex string) (*imageMetadata, error) { diff --git a/lib/images/storage_test.go b/lib/images/storage_test.go index b34c47d54..d7c8fd994 100644 --- a/lib/images/storage_test.go +++ b/lib/images/storage_test.go @@ -7,7 +7,6 @@ import ( "time" "github.com/kernel/hypeman/lib/paths" - "github.com/kernel/hypeman/lib/tags" "github.com/stretchr/testify/require" ) @@ -162,79 +161,6 @@ func TestContentLayoutResolvesDiskByDigest(t *testing.T) { 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)) - - mgr := &manager{paths: p} - image, err := mgr.GetImage(nil, second) - require.NoError(t, err) - require.Equal(t, tags.Tags{"team": "two"}, image.Tags) - - 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 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" @@ -266,12 +192,46 @@ func TestListAllMetadataContentLayout(t *testing.T) { }, names) } +func TestImageMetadataToImage_ClonesMetadata(t *testing.T) { + createdAt := time.Now().UTC().Truncate(time.Second) + source := &imageMetadata{ + Name: "docker.io/library/alpine:latest", + Digest: "sha256:abc", + Status: StatusReady, + Tags: map[string]string{"team": "backend", "env": "staging"}, + SizeBytes: 123, + CreatedAt: createdAt, + } + + img := source.toImageFor(source.Name) + 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) + + source.Tags["team"] = "mutated" + require.Equal(t, "backend", img.Tags["team"]) +} + +func TestImageMetadataToImage_EmptyMetadataOmitted(t *testing.T) { + img := (&imageMetadata{ + Name: "docker.io/library/alpine:latest", + Digest: "sha256:abc", + Status: StatusPending, + CreatedAt: time.Now().UTC(), + }).toImageFor("docker.io/library/alpine:latest") + + require.Nil(t, img.Tags) +} + func TestPromoteLegacyImagesMovesContentAndTags(t *testing.T) { p := paths.New(t.TempDir()) repository := "docker.io/library/alpine" tag := "latest" digest := "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc" + // Ready legacy image with a legacy tag symlink. legacyDir := p.ImageDigestDir(repository, digest) require.NoError(t, os.MkdirAll(legacyDir, 0o755)) meta := &imageMetadata{ @@ -289,6 +249,7 @@ func TestPromoteLegacyImagesMovesContentAndTags(t *testing.T) { promoteLegacyImages(p) + // Content exists with the same bytes and is ready. contentMeta, err := readContentMetadata(p, digest) require.NoError(t, err) require.Equal(t, StatusReady, contentMeta.Status) @@ -298,6 +259,7 @@ func TestPromoteLegacyImagesMovesContentAndTags(t *testing.T) { require.FileExists(t, p.ImageDigestPath(repository, digest)) + // The tag now resolves through the shared content layout. resolved, err := resolveTag(p, repository, tag) require.NoError(t, err) require.Equal(t, digest, resolved) diff --git a/lib/images/tag.go b/lib/images/tag.go deleted file mode 100644 index 181176952..000000000 --- a/lib/images/tag.go +++ /dev/null @@ -1,106 +0,0 @@ -package images - -import ( - "context" - "errors" - "fmt" - "log/slog" -) - -// TagImage creates a ready-image tag without pulling or converting content. -// Cross-repository tags promote legacy content into the shared layout. A -// failed call leaves no side effects: the target tag's generation is only -// bumped after the new tag is on disk, so pending pulls that claimed the -// target tag keep their claim. When the target previously pointed at -// different content, that digest is collected after the new tag is live; -// cleanup failures are logged and do not fail the call, since the tag is -// already installed at that point. -// -// Promotion deliberately runs before the symlink install, so a symlink -// failure after a cross-repo promotion leaves the content promoted with no -// target tag. That state is gc-consistent (unreferenced content is -// collected) and retry is idempotent (promoteImageToContent short-circuits -// on ready content), which beats rolling back a completed promotion. -func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, error) { - sourceRef, targetRef, err := parseTagReferences(source, target) - if err != nil { - return nil, err - } - - m.createMu.Lock() - defer m.createMu.Unlock() - - digestHex, meta, err := m.readyTagImage(sourceRef) - if err != nil { - return nil, err - } - - // A dangling or malformed target symlink is treated like a missing tag so - // the retag self-heals; createTagSymlink replaces the link either way. - previousDigest, err := resolveTag(m.paths, targetRef.Repository(), targetRef.Tag()) - if err != nil && !errors.Is(err, ErrNotFound) && !errors.Is(err, errInvalidSymlinkTarget) { - return nil, fmt.Errorf("resolve existing target tag: %w", err) - } - if sourceRef.Repository() != targetRef.Repository() { - if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta); err != nil { - return nil, fmt.Errorf("promote image to content: %w", err) - } - } - if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { - return nil, fmt.Errorf("create image tag: %w", err) - } - // Unlike updateExistingReference, which bumps the generation before - // installing the symlink, the bump happens after the install here so a - // failed tag call cannot invalidate a pending pull's claim on the target - // tag (pinned by TestTagImageFailureLeavesNoSideEffects). - m.nextTagGeneration(targetRef.Repository(), targetRef.Tag()) - m.cleanupReplacedTag(targetRef, previousDigest, digestHex) - - return meta.toImageFor(targetRef.String()), nil -} - -func (m *manager) cleanupReplacedTag(ref *NormalizedRef, previousDigest, digestHex string) { - if previousDigest == "" || previousDigest == digestHex { - return - } - // Sibling tags in this repository may still reference the previous - // digest; only collect when this was the last reference. - count, err := countTagsForDigest(m.paths, ref.Repository(), previousDigest) - if err != nil { - slog.Warn("failed to count tags for replaced image", "repository", ref.Repository(), "digest", previousDigest, "error", err) - return - } - if count > 0 { - return - } - if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), previousDigest, true); err != nil { - slog.Warn("failed to collect replaced image content", "repository", ref.Repository(), "digest", previousDigest, "error", err) - } - m.refreshDiskUsageTotals() -} - -func parseTagReferences(source, target string) (*NormalizedRef, *NormalizedRef, error) { - sourceRef, err := ParseNormalizedRef(source) - if err != nil { - return nil, nil, fmt.Errorf("%w: invalid source reference: %s", ErrInvalidName, err) - } - targetRef, err := ParseNormalizedRef(target) - if err != nil { - return nil, nil, fmt.Errorf("%w: invalid target reference: %s", ErrInvalidName, err) - } - if targetRef.IsDigest() { - return nil, nil, fmt.Errorf("%w: target must be a tag reference, not a digest", ErrInvalidName) - } - return sourceRef, targetRef, nil -} - -func (m *manager) readyTagImage(ref *NormalizedRef) (string, *imageMetadata, error) { - digestHex, meta, err := resolveRefMetadata(m.paths, ref) - if err != nil { - return "", nil, err - } - if meta.Status != StatusReady { - return "", nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) - } - return digestHex, meta, nil -} diff --git a/lib/images/tag_manager.go b/lib/images/tag_manager.go new file mode 100644 index 000000000..94f5954c1 --- /dev/null +++ b/lib/images/tag_manager.go @@ -0,0 +1,261 @@ +package images + +import ( + "context" + "errors" + "fmt" + "log/slog" + "os" + "strings" +) + +func tagGenerationKey(repository, tag string) string { + return repository + ":" + tag +} + +// nextTagGeneration bumps the mutation generation for a tag and returns the +// new value. Pulls record the value with their metadata and only repoint the +// tag on completion when no later operation mutated it. +func (m *manager) nextTagGeneration(repository, tag string) uint64 { + if m.tagGenerations == nil { + m.tagGenerations = make(map[string]uint64) + } + key := tagGenerationKey(repository, tag) + m.tagGenerations[key]++ + return m.tagGenerations[key] +} + +// revertTagGeneration undoes a nextTagGeneration bump for an operation that +// failed before mutating the tag, so an in-flight pull of the tag is not +// permanently blocked from repointing it. +func (m *manager) revertTagGeneration(repository, tag string) { + key := tagGenerationKey(repository, tag) + if m.tagGenerations[key] <= 1 { + delete(m.tagGenerations, key) + return + } + m.tagGenerations[key]-- +} + +// releaseTagGeneration drops a failed pull's claim on its tag so an older +// in-flight pull of the same tag can still repoint it. +func (m *manager) releaseTagGeneration(repository, tag string, generation uint64) { + if m.tagGenerations[tagGenerationKey(repository, tag)] == generation { + m.revertTagGeneration(repository, tag) + } +} + +// pruneTagGenerations drops entries for tags that no longer resolve, keeping +// the map bounded over the process lifetime. Deleting a tag still suppresses +// an in-flight pull's repoint: the pull's recorded generation can no longer +// match the missing entry. +func (m *manager) pruneTagGenerations() { + for key := range m.tagGenerations { + // tagGenerationKey is repo+":"+tag and tags never contain ":", but + // repositories may (host:port), so split on the last colon. + colon := strings.LastIndexByte(key, ':') + if colon < 0 { + continue + } + if _, err := resolveTag(m.paths, key[:colon], key[colon+1:]); err != nil { + delete(m.tagGenerations, key) + } + } +} + +// restoreTagState re-seeds the tag indexes from recovered metadata. metas must +// be sorted oldest first so the newest requested pull wins. +func (m *manager) restoreTagState(metas []*imageMetadata) { + if m.tagGenerations == nil { + m.tagGenerations = make(map[string]uint64) + } + if m.requestedTags == nil { + m.requestedTags = make(map[string]string) + } + for _, meta := range metas { + if meta.RequestedTag == "" || meta.Digest == "" { + continue + } + ref, err := ParseNormalizedRef(meta.Name) + if err != nil { + continue + } + key := tagGenerationKey(ref.Repository(), meta.RequestedTag) + if meta.TagGeneration > m.tagGenerations[key] { + m.tagGenerations[key] = meta.TagGeneration + } + m.requestedTags[key] = strings.TrimPrefix(meta.Digest, "sha256:") + } +} + +// claimTagForStatus repoints ref's tag at the digest when the image is ready, +// or records a pending tag when the build is still in flight. It is a no-op +// for digest-only references. +func (m *manager) claimTagForStatus(meta *imageMetadata, ref *ResolvedRef) error { + if ref.Tag() == "" { + return nil + } + if meta.Status == StatusReady { + return m.claimReadyTag(ref.Repository(), ref.Tag(), ref.DigestHex()) + } + return ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) +} + +// claimReadyTag repoints an existing tag at a ready digest, last pull wins. +func (m *manager) claimReadyTag(repository, tag, digestHex string) error { + m.nextTagGeneration(repository, tag) + if err := createTagSymlink(m.paths, repository, tag, digestHex); err != nil { + m.revertTagGeneration(repository, tag) + return err + } + return nil +} + +// trackRequestedTag records the digest of the newest pull requested for a tag +// so readiness waits can find it without walking the metadata tree. +func (m *manager) trackRequestedTag(repository, tag, digestHex string) { + if m.requestedTags == nil { + m.requestedTags = make(map[string]string) + } + m.requestedTags[tagGenerationKey(repository, tag)] = digestHex +} + +// requestedTagImage returns the newest image requested for a tag, so a +// readiness wait tracks the latest pull rather than the digest the tag +// currently points at. +func (m *manager) requestedTagImage(ref *NormalizedRef) *Image { + m.createMu.Lock() + digestHex, ok := m.requestedTags[tagGenerationKey(ref.Repository(), ref.Tag())] + m.createMu.Unlock() + if !ok { + return nil + } + meta, err := readMetadata(m.paths, ref.Repository(), digestHex) + if err != nil { + return nil + } + return meta.toImageFor(ref.String()) +} + +// claimRequestedTag repoints the pull's requested tag at the finished digest. +// The tag is only repointed when its generation still matches the pull's +// recorded claim, so a later mutation of the tag wins over the in-flight +// pull. A recovered build whose metadata predates requested-tag tracking has +// no recorded tag: fall back to the reference's own tag and recreate a +// missing symlink, matching the pre-tracking behavior. Otherwise an +// already-missing tag stays missing so a concurrent delete wins. +func (m *manager) claimRequestedTag(ref *ResolvedRef, meta *imageMetadata) { + requestedTag := meta.RequestedTag + allowMissing := requestedTag == "" + if requestedTag == "" { + requestedTag = ref.Tag() + } + if requestedTag == "" || m.tagGenerations[tagGenerationKey(ref.Repository(), requestedTag)] != meta.TagGeneration { + return + } + current, err := resolveTag(m.paths, ref.Repository(), requestedTag) + if err != nil { + if !allowMissing || !errors.Is(err, ErrNotFound) { + return + } + } else if current != ref.DigestHex() && current != meta.PreviousTagDigest { + return + } + if err := createTagSymlink(m.paths, ref.Repository(), requestedTag, ref.DigestHex()); err != nil { + fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) + } +} + +// TagImage creates a ready-image tag without pulling or converting content. +// Cross-repository tags promote legacy content into the shared layout. A +// failed call leaves no side effects: the target tag's generation is only +// bumped after the new tag is on disk, so pending pulls that claimed the +// target tag keep their claim. When the target previously pointed at +// different content, that digest is collected after the new tag is live; +// cleanup failures are logged and do not fail the call, since the tag is +// already installed at that point. +// +// Promotion deliberately runs before the symlink install, so a symlink +// failure after a cross-repo promotion leaves the content promoted with no +// target tag. That state is gc-consistent (unreferenced content is +// collected) and retry is idempotent (promoteImageToContent short-circuits +// on ready content), which beats rolling back a completed promotion. +func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, error) { + sourceRef, targetRef, err := parseTagReferences(source, target) + if err != nil { + return nil, err + } + + m.createMu.Lock() + defer m.createMu.Unlock() + + // A dangling or malformed target symlink is treated like a missing tag so + // the retag self-heals; createTagSymlink replaces the link either way. + previousDigest, err := resolveTag(m.paths, targetRef.Repository(), targetRef.Tag()) + if err != nil && !errors.Is(err, ErrNotFound) && !errors.Is(err, errInvalidSymlinkTarget) { + return nil, fmt.Errorf("resolve existing target tag: %w", err) + } + + digestHex, meta, err := m.readyTagImage(sourceRef) + if err != nil { + return nil, err + } + if sourceRef.Repository() != targetRef.Repository() { + if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta); err != nil { + return nil, fmt.Errorf("promote image to content: %w", err) + } + } + if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { + return nil, fmt.Errorf("create image tag: %w", err) + } + // Unlike updateExistingReference, which bumps the generation before + // installing the symlink, the bump happens after the install here so a + // failed tag call cannot invalidate a pending pull's claim on the target + // tag (pinned by TestTagImageFailureLeavesNoSideEffects). + m.nextTagGeneration(targetRef.Repository(), targetRef.Tag()) + m.cleanupReplacedTag(targetRef, previousDigest, digestHex) + + return meta.toImageFor(targetRef.String()), nil +} + +func (m *manager) cleanupReplacedTag(ref *NormalizedRef, previousDigest, digestHex string) { + if previousDigest != "" && previousDigest != digestHex { + // Sibling tags in this repository may still reference the previous + // digest; only collect when this was the last reference. + count, err := countTagsForDigest(m.paths, ref.Repository(), previousDigest) + if err != nil { + slog.Warn("failed to count tags for replaced image", "repository", ref.Repository(), "digest", previousDigest, "error", err) + } else if count == 0 { + if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), previousDigest, true); err != nil { + slog.Warn("failed to collect replaced image content", "repository", ref.Repository(), "digest", previousDigest, "error", err) + } + m.refreshDiskUsageTotals() + } + } +} + +func parseTagReferences(source, target string) (*NormalizedRef, *NormalizedRef, error) { + sourceRef, err := ParseNormalizedRef(source) + if err != nil { + return nil, nil, fmt.Errorf("%w: invalid source reference: %s", ErrInvalidName, err) + } + targetRef, err := ParseNormalizedRef(target) + if err != nil { + return nil, nil, fmt.Errorf("%w: invalid target reference: %s", ErrInvalidName, err) + } + if targetRef.IsDigest() { + return nil, nil, fmt.Errorf("%w: target must be a tag reference, not a digest", ErrInvalidName) + } + return sourceRef, targetRef, nil +} + +func (m *manager) readyTagImage(ref *NormalizedRef) (string, *imageMetadata, error) { + digestHex, meta, err := resolveRefMetadata(m.paths, ref) + if err != nil { + return "", nil, err + } + if meta.Status != StatusReady { + return "", nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) + } + return digestHex, meta, nil +} diff --git a/lib/paths/paths.go b/lib/paths/paths.go index 814dc1432..bc9403cd6 100644 --- a/lib/paths/paths.go +++ b/lib/paths/paths.go @@ -161,6 +161,12 @@ func (p *Paths) ImageContentMetadata(digestHex string) string { return filepath.Join(p.ImageContentDir(digestHex), "metadata.json") } +// ImageContentManifestModel returns the path to the persisted OCI manifest +// model (layer descriptors, config, platform) for content-addressed image data. +func (p *Paths) ImageContentManifestModel(digestHex string) string { + return filepath.Join(p.ImageContentDir(digestHex), "manifest.json") +} + // ImageRepositoriesDir returns the root directory for repository tag references. func (p *Paths) ImageRepositoriesDir() string { return filepath.Join(p.dataDir, "images", "repositories")