diff --git a/lib/diskutilization/diskutilization.go b/lib/diskutilization/diskutilization.go index 27fef3c81..a5f55d677 100644 --- a/lib/diskutilization/diskutilization.go +++ b/lib/diskutilization/diskutilization.go @@ -4,6 +4,7 @@ import ( "io/fs" "os" "path/filepath" + "strings" "syscall" "github.com/kernel/hypeman/lib/paths" @@ -56,7 +57,7 @@ func Collect(p *paths.Paths) (Breakdown, error) { return false } name := entry.Name() - return name == "rootfs.erofs" || name == "rootfs.ext4" + return name == "rootfs.erofs" || name == "rootfs.ext4" || strings.HasPrefix(name, "layer.") }) if err != nil { return Breakdown{}, err @@ -178,6 +179,7 @@ func sumDirectChildFileAllocatedBytes(root string, childFile string) (int64, err func sumMatchingFilesAllocatedBytes(root string, match func(path string, entry fs.DirEntry) bool) (int64, error) { var total int64 + seen := make(map[fileIdentity]struct{}) err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { if err != nil { if os.IsNotExist(err) { @@ -185,9 +187,24 @@ func sumMatchingFilesAllocatedBytes(root string, match func(path string, entry f } return err } - if match(path, entry) { - total += allocatedBytesForPath(path) + if !match(path, entry) { + return nil } + info, statErr := os.Lstat(path) + if statErr != nil { + if os.IsNotExist(statErr) { + return nil + } + return statErr + } + if stat, ok := info.Sys().(*syscall.Stat_t); ok { + identity := fileIdentity{dev: uint64(stat.Dev), ino: uint64(stat.Ino)} + if _, exists := seen[identity]; exists { + return nil + } + seen[identity] = struct{}{} + } + total += allocatedBytesForPath(path) return nil }) if err != nil { @@ -244,6 +261,11 @@ func sumSnapshotTreeAllocatedBytes(root string, sharedExtents *sharedExtentTrack return privateTotal, sharedTotal, nil } +type fileIdentity struct { + dev uint64 + ino uint64 +} + func allocatedBytesForPath(path string) int64 { info, err := os.Lstat(path) if err != nil { diff --git a/lib/diskutilization/diskutilization_test.go b/lib/diskutilization/diskutilization_test.go index 50a787017..ff854438d 100644 --- a/lib/diskutilization/diskutilization_test.go +++ b/lib/diskutilization/diskutilization_test.go @@ -98,6 +98,22 @@ func TestCollect_UsesAllocatedBytesAndClassifiesSnapshots(t *testing.T) { require.Equal(t, otherTotal, utilization.SnapshotOther) } +func TestCollect_DeduplicatesHardLinkedImagesAndCountsLayers(t *testing.T) { + p := paths.New(t.TempDir()) + imagePath := filepath.Join(p.ImagesDir(), "repo", "digest", "rootfs.erofs") + require.NoError(t, createSparseTestFile(imagePath, 8192, []sparseWrite{{offset: 0, data: []byte("image")}})) + aliasPath := filepath.Join(p.ImagesDir(), "content", "digest", "rootfs.erofs") + require.NoError(t, os.MkdirAll(filepath.Dir(aliasPath), 0755)) + require.NoError(t, os.Link(imagePath, aliasPath)) + + layerPath := filepath.Join(p.ImageLayersDir(), "layer-digest", "layer.erofs") + require.NoError(t, createSparseTestFile(layerPath, 8192, []sparseWrite{{offset: 0, data: []byte("layer")}})) + + utilization, err := Collect(p) + require.NoError(t, err) + require.Equal(t, allocatedBytesForPath(imagePath)+allocatedBytesForPath(layerPath), utilization.Images) +} + func createSparseTestFile(path string, size int64, writes []sparseWrite) error { if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil { return err diff --git a/lib/images/compose.go b/lib/images/compose.go index a5e0d7a2a..718fadd50 100644 --- a/lib/images/compose.go +++ b/lib/images/compose.go @@ -4,11 +4,11 @@ import ( "fmt" "os" "path/filepath" - "strings" ) -// validateModelPairing validates the persisted model before composition so -// every manifest layer has a corresponding, verified config diff ID. +// validateModelPairing mirrors validateConfigFileForUnpack: the image config +// must carry one diff id per manifest layer so composition never indexes past +// the end of the pairing. func validateModelPairing(layoutTag string, model *imageManifestModel) error { if err := validateManifestModel(layoutTag, model); err != nil { return fmt.Errorf("unpack rootfs: %w", err) @@ -17,51 +17,79 @@ func validateModelPairing(layoutTag string, model *imageManifestModel) error { } // composeRootfs merges an image's layers into dest in manifest order, reading -// each layer blob from the shared OCI cache. Whiteout and opaque-directory -// markers are interpreted as each layer is applied. +// each layer blob from the shared OCI cache. The result is one complete rootfs +// tree that is exported to a single disk, matching the guest's contract: one +// read-only lower filesystem and one writable overlay upper. Whiteout and +// opaque-directory markers are interpreted as each layer is applied instead of +// being left in the tree, so the composed rootfs never relies on tar-level +// whiteouts composing on overlayfs. func (c *ociClient) composeRootfs(dest string, layers []layerDescriptor) error { + trees, err := c.composeRootfsWithLayerTrees(dest, layers) + if err != nil { + return err + } + cleanupLayerTrees(trees) + return nil +} + +func (c *ociClient) composeRootfsWithLayerTrees(dest string, layers []layerDescriptor) (map[string]layerTree, error) { if len(layers) == 0 { - return fmt.Errorf("image has no layers") + return nil, fmt.Errorf("image has no layers") } if err := os.MkdirAll(dest, 0755); err != nil { - return fmt.Errorf("create compose directory: %w", err) + return nil, fmt.Errorf("create compose directory: %w", err) } + trees := make(map[string]layerTree, len(layers)) for i, desc := range layers { - if err := c.applyLayerToDir(dest, desc); err != nil { - return fmt.Errorf("apply layer %d (%s): %w", i, desc.Digest, err) + if tree, ok := trees[desc.Digest]; ok { + if err := applyLayerTree(tree.path, dest); err != nil { + cleanupLayerTrees(trees) + return nil, fmt.Errorf("apply layer %d (%s): %w", i, desc.Digest, err) + } + continue + } + tree, err := c.extractLayerTree(desc) + if err != nil { + cleanupLayerTrees(trees) + return nil, fmt.Errorf("extract layer %d (%s): %w", i, desc.Digest, err) } + if err := applyLayerTree(tree.path, dest); err != nil { + cleanupLayerTrees(trees) + cleanupLayerTree(tree) + return nil, fmt.Errorf("apply layer %d (%s): %w", i, desc.Digest, err) + } + trees[desc.Digest] = tree } - return nil + return trees, nil } -func (c *ociClient) applyLayerToDir(dest string, desc layerDescriptor) error { - layerHex := strings.TrimPrefix(desc.Digest, "sha256:") - if layerHex == "" || strings.Contains(layerHex, "/") || layerHex == "." || strings.Contains(layerHex, "..") { - return fmt.Errorf("invalid layer digest: %s", desc.Digest) +// extractLayerTree extracts one layer into a private staging directory so it +// can be applied to a composed rootfs and later materialized as an artifact. +func (c *ociClient) extractLayerTree(desc layerDescriptor) (layerTree, error) { + layerHex, err := layerDigestHex(desc) + if err != nil { + return layerTree{}, err } blobPath := filepath.Join(c.cacheDir, "blobs", "sha256", layerHex) if _, err := os.Stat(blobPath); err != nil { if os.IsNotExist(err) { - return fmt.Errorf("layer blob missing from oci cache: %s", desc.Digest) + return layerTree{}, fmt.Errorf("layer blob missing from oci cache: %s", desc.Digest) } - return fmt.Errorf("stat layer blob: %w", err) + return layerTree{}, fmt.Errorf("stat layer blob: %w", err) } layerDir, err := os.MkdirTemp("", "hypeman-layer-*") if err != nil { - return fmt.Errorf("create layer staging directory: %w", err) + return layerTree{}, fmt.Errorf("create layer staging directory: %w", err) } - defer os.RemoveAll(layerDir) - stats, err := unpackLayerBlob(blobPath, desc.MediaType, layerDir) if err != nil { - return err + cleanupLayerTree(layerTree{path: layerDir}) + return layerTree{}, err } if desc.DiffID != "" && stats.diffID != desc.DiffID { - return fmt.Errorf("layer %s diff id mismatch: got %s, want %s", desc.Digest, stats.diffID, desc.DiffID) - } - if err := applyLayerTree(layerDir, dest); err != nil { - return fmt.Errorf("apply layer tree: %w", err) + cleanupLayerTree(layerTree{path: layerDir}) + return layerTree{}, fmt.Errorf("layer %s diff id mismatch: got %s, want %s", desc.Digest, stats.diffID, desc.DiffID) } - return nil + return layerTree{path: layerDir, stats: stats}, nil } diff --git a/lib/images/credentials_test.go b/lib/images/credentials_test.go index 7f680b980..87a72c39b 100644 --- a/lib/images/credentials_test.go +++ b/lib/images/credentials_test.go @@ -83,10 +83,8 @@ func TestCreateImageRequestCredentialsAreNotPersisted(t *testing.T) { } func TestInflightPullRejectsDifferentCredentials(t *testing.T) { - m := &manager{ - inflightPulls: make(map[string]*inflightImagePull), - borrowedCredentialsTimeout: time.Minute, - } + m := newTestManager(nil) + m.borrowedCredentialsTimeout = time.Minute const digest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" credentials := &authn.AuthConfig{Username: "AWS", Password: "token-a"} inflight := m.registerInflightPull(digest, credentials) @@ -98,10 +96,8 @@ func TestInflightPullRejectsDifferentCredentials(t *testing.T) { } func TestBorrowedCredentialsExpireWhileQueued(t *testing.T) { - m := &manager{ - inflightPulls: make(map[string]*inflightImagePull), - borrowedCredentialsTimeout: time.Millisecond, - } + m := newTestManager(nil) + m.borrowedCredentialsTimeout = time.Millisecond const digest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" inflight := m.registerInflightPull(digest, &authn.AuthConfig{Username: "AWS", Password: "secret"}) defer m.releaseInflightPull(digest, inflight)() @@ -119,7 +115,7 @@ func TestBorrowedCredentialsExpireWhileQueued(t *testing.T) { } func TestBorrowedAuthRejectsReplacedInflightPull(t *testing.T) { - m := &manager{inflightPulls: make(map[string]*inflightImagePull)} + m := newTestManager(nil) const digest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" first := m.registerInflightPull(digest, &authn.AuthConfig{Username: "first"}) second := m.registerInflightPull(digest, &authn.AuthConfig{Username: "second"}) @@ -174,12 +170,9 @@ func TestRecoverInterruptedCredentialedPullFailsForFreshRetry(t *testing.T) { p := paths.New(t.TempDir()) client, err := newOCIClient(p.SystemOCICache()) require.NoError(t, err) - m := &manager{ - paths: p, - ociClient: client, - queue: queue.New(1), - readySubscribers: make(map[string][]chan StatusEvent), - } + m := newTestManager(p) + m.ociClient = client + m.queue = queue.New(1) const repository = "registry.example/private/image" const digest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" diff --git a/lib/images/disk_usage.go b/lib/images/disk_usage.go index b65dba1ac..c1cffbb2c 100644 --- a/lib/images/disk_usage.go +++ b/lib/images/disk_usage.go @@ -108,57 +108,80 @@ func totalOCICacheBlobBytesFromFilesystem(blobDir string) (int64, error) { return total, nil } -func (m *manager) getDiskUsageTotals() (int64, int64, error) { +// diskUsageTotals caches the size of each disk component the manager tracks. +type diskUsageTotals struct { + readyImageBytes int64 + layerBytes int64 + ociCacheBytes int64 +} + +func (m *manager) getDiskUsageTotals() (diskUsageTotals, error) { m.diskUsageMu.RLock() if m.diskUsageLoaded { - readyImageBytes := m.readyImageBytes - ociCacheBytes := m.ociCacheBytes + totals := diskUsageTotals{ + readyImageBytes: m.readyImageBytes, + layerBytes: m.layerBytes, + ociCacheBytes: m.ociCacheBytes, + } m.diskUsageMu.RUnlock() - return readyImageBytes, ociCacheBytes, nil + return totals, nil } m.diskUsageMu.RUnlock() - readyImageBytes, ociCacheBytes, err := m.computeDiskUsageTotals() + computed, err := m.computeDiskUsageTotals() if err != nil { - return 0, 0, err + return diskUsageTotals{}, err } m.diskUsageMu.Lock() if !m.diskUsageLoaded { - m.readyImageBytes = readyImageBytes - m.ociCacheBytes = ociCacheBytes + m.readyImageBytes = computed.readyImageBytes + m.layerBytes = computed.layerBytes + m.ociCacheBytes = computed.ociCacheBytes m.diskUsageLoaded = true } - readyImageBytes = m.readyImageBytes - ociCacheBytes = m.ociCacheBytes + totals := diskUsageTotals{ + readyImageBytes: m.readyImageBytes, + layerBytes: m.layerBytes, + ociCacheBytes: m.ociCacheBytes, + } m.diskUsageMu.Unlock() - return readyImageBytes, ociCacheBytes, nil + return totals, nil } func (m *manager) refreshDiskUsageTotals() { - readyImageBytes, ociCacheBytes, err := m.computeDiskUsageTotals() + computed, err := m.computeDiskUsageTotals() if err != nil { return } m.diskUsageMu.Lock() - m.readyImageBytes = readyImageBytes - m.ociCacheBytes = ociCacheBytes + m.readyImageBytes = computed.readyImageBytes + m.layerBytes = computed.layerBytes + m.ociCacheBytes = computed.ociCacheBytes m.diskUsageLoaded = true m.diskUsageMu.Unlock() } -func (m *manager) computeDiskUsageTotals() (int64, int64, error) { +func (m *manager) computeDiskUsageTotals() (diskUsageTotals, error) { readyImageBytes, err := totalReadyImageBytesFromMetadata(m.paths.ImagesDir()) if err != nil { - return 0, 0, err + return diskUsageTotals{}, err + } + layerBytes, err := totalLayerArtifactBytes(m.paths.ImageLayersDir()) + if err != nil { + return diskUsageTotals{}, err } ociCacheBytes, err := totalOCICacheBlobBytesFromFilesystem(m.paths.OCICacheBlobDir()) if err != nil { - return 0, 0, err + return diskUsageTotals{}, err } - return readyImageBytes, ociCacheBytes, nil + return diskUsageTotals{ + readyImageBytes: readyImageBytes, + layerBytes: layerBytes, + ociCacheBytes: ociCacheBytes, + }, nil } func totalRootfsBytesInDigestDir(digestDir string) (int64, error) { diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index 41d7201e8..7193baeac 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -110,16 +110,11 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { return &record, nil } -// materializeLayerArtifact ensures a layer has a materialized artifact keyed -// by its blob digest, building it from the shared OCI cache blob when absent. -// The layer is unpacked into an isolated temp directory, converted to erofs, -// and installed atomically; an interrupted build leaves only temp files that -// the next attempt replaces. func (a *layerArtifact) validate() error { if a.SchemaVersion != layerRecordSchemaVersion { return fmt.Errorf("unsupported schema version: %d", a.SchemaVersion) } - if a.Digest == "" || a.Format != layerFormatErofs && a.Format != layerFormatExt4 { + if a.Digest == "" || (a.Format != layerFormatErofs && a.Format != layerFormatExt4) { return fmt.Errorf("invalid digest or format") } if a.SizeBytes < 0 || a.UnpackedBytes < 0 || a.Entries < 0 { @@ -134,55 +129,48 @@ func (a *layerArtifact) validate() error { return nil } -func (m *manager) materializeLayerArtifact(desc layerDescriptor) (*layerArtifact, error) { - layerHex := strings.TrimPrefix(desc.Digest, "sha256:") - if err := paths.ValidatePathComponent(layerHex); err != nil { - return nil, fmt.Errorf("invalid layer digest: %s", desc.Digest) - } - - if record, err := readLayerRecord(m.paths, layerHex); err != nil { +// materializeLayerArtifact ensures a layer has a materialized artifact keyed +// by its blob digest, converting the layer's staging tree from the pull and +// installing it atomically; an interrupted build leaves only temp files that +// the next attempt replaces. +func (m *manager) materializeLayerArtifact(desc layerDescriptor, tree layerTree) (*layerArtifact, error) { + layerHex, err := layerDigestHex(desc) + if err != nil { return nil, err - } else if record != nil && record.matches(desc) { - if _, statErr := os.Stat(layerArtifactPath(m.paths, layerHex)); statErr == nil { - return record, nil - } - // Record without artifact: rebuild below. } - - blobPath := m.paths.OCICacheBlob(layerHex) - if _, err := os.Stat(blobPath); err != nil { - if os.IsNotExist(err) { - return nil, fmt.Errorf("layer blob missing from oci cache: %s", desc.Digest) - } - return nil, fmt.Errorf("stat layer blob: %w", err) + if desc.DiffID != "" && tree.stats.diffID != desc.DiffID { + return nil, fmt.Errorf("layer %s diff id mismatch: got %s, want %s", desc.Digest, tree.stats.diffID, desc.DiffID) } - layerDir := m.paths.ImageLayerDir(layerHex) - if err := os.MkdirAll(layerDir, 0755); err != nil { - return nil, fmt.Errorf("create layer directory: %w", err) + lock := m.layerDigestLock(layerHex) + lock.Lock() + defer lock.Unlock() + if record, err := m.existingLayerArtifact(layerHex, desc); record != nil || err != nil { + return record, err } - unpackDir, err := os.MkdirTemp(layerDir, ".unpack-*") - if err != nil { - return nil, fmt.Errorf("create unpack directory: %w", err) + return m.installLayerArtifact(desc, layerHex, tree.path, tree.stats) +} + +func layerDigestHex(desc layerDescriptor) (string, error) { + layerHex := strings.TrimPrefix(desc.Digest, "sha256:") + if err := paths.ValidatePathComponent(layerHex); err != nil { + return "", fmt.Errorf("invalid layer digest: %s", desc.Digest) } - defer os.RemoveAll(unpackDir) + return layerHex, nil +} - stats, err := unpackLayerBlob(blobPath, desc.MediaType, unpackDir) +func (m *manager) existingLayerArtifact(layerHex string, desc layerDescriptor) (*layerArtifact, error) { + record, err := readLayerRecord(m.paths, layerHex) if err != nil { - return nil, fmt.Errorf("unpack layer %s: %w", desc.Digest, err) + return nil, err } - if desc.DiffID != "" && stats.diffID != desc.DiffID { - return nil, fmt.Errorf("layer %s diff id mismatch: got %s, want %s", desc.Digest, stats.diffID, desc.DiffID) + if record == nil || !record.matches(desc) { + return nil, nil } - - return m.installLayerArtifact(desc, layerHex, unpackDir, stats) -} - -func artifactOptions(format string) layerArtifactOptions { - if format == layerFormatErofs { - return layerArtifactOptions{Compression: "lz4"} + if _, err := os.Stat(layerArtifactPath(m.paths, layerHex)); err != nil { + return nil, nil } - return layerArtifactOptions{} + return record, nil } func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir string, stats *unpackStats) (*layerArtifact, error) { @@ -227,6 +215,18 @@ func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir return record, nil } +func artifactOptions(format string) layerArtifactOptions { + if format == layerFormatErofs { + return layerArtifactOptions{Compression: "lz4"} + } + return layerArtifactOptions{} +} + +type layerTree struct { + path string + stats *unpackStats +} + type unpackStats struct { entries int unpackedBytes int64 @@ -234,6 +234,18 @@ type unpackStats struct { whiteouts []whiteoutRecord } +func cleanupLayerTree(tree layerTree) { + if tree.path != "" { + _ = os.RemoveAll(tree.path) + } +} + +func cleanupLayerTrees(trees map[string]layerTree) { + for _, tree := range trees { + cleanupLayerTree(tree) + } +} + // unpackLayerBlob extracts one compressed layer blob into dest, preserving // whiteout marker files and recording them. Paths are confined to dest. func unpackLayerBlob(blobPath, mediaType, dest string) (*unpackStats, error) { diff --git a/lib/images/layer_artifact_test.go b/lib/images/layer_artifact_test.go index 79d00e01b..9f8648bd6 100644 --- a/lib/images/layer_artifact_test.go +++ b/lib/images/layer_artifact_test.go @@ -55,9 +55,15 @@ func TestMaterializeLayerArtifact(t *testing.T) { writeLayerTestLayout(t, p, img) desc := layerDescFromImage(t, img, 0) - m := &manager{paths: p} + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) + m := newTestManager(p) + + tree, err := client.extractLayerTree(desc) + require.NoError(t, err) + defer cleanupLayerTree(tree) - record, err := m.materializeLayerArtifact(desc) + record, err := m.materializeLayerArtifact(desc, tree) require.NoError(t, err) require.Equal(t, desc.Digest, record.Digest) require.Equal(t, desc.DiffID, record.DiffID) @@ -73,7 +79,7 @@ func TestMaterializeLayerArtifact(t *testing.T) { // A second materialization reuses the existing artifact. artifactInfo, err := os.Stat(p.ImageLayerArtifact(layerHex)) require.NoError(t, err) - reused, err := m.materializeLayerArtifact(desc) + reused, err := m.materializeLayerArtifact(desc, tree) require.NoError(t, err) require.True(t, record.CreatedAt.Equal(reused.CreatedAt), "reuse must return the stored record") artifactInfoAfter, err := os.Stat(p.ImageLayerArtifact(layerHex)) @@ -81,11 +87,12 @@ func TestMaterializeLayerArtifact(t *testing.T) { require.Equal(t, artifactInfo.ModTime(), artifactInfoAfter.ModTime(), "reuse must not rebuild") } -func TestMaterializeLayerArtifactMissingBlob(t *testing.T) { +func TestExtractLayerTreeMissingBlob(t *testing.T) { p := paths.New(t.TempDir()) - m := &manager{paths: p} + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) - _, err := m.materializeLayerArtifact(layerDescriptor{ + _, err = client.extractLayerTree(layerDescriptor{ Digest: "sha256:abababababababababababababababababababababababababababababababab", MediaType: "application/vnd.oci.image.layer.v1.tar+gzip", }) @@ -136,9 +143,15 @@ func TestMaterializeLayerRecordsWhiteouts(t *testing.T) { writeLayerTestLayout(t, p, img) desc := layerDescFromImage(t, img, 0) - m := &manager{paths: p} + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) + m := newTestManager(p) + + tree, err := client.extractLayerTree(desc) + require.NoError(t, err) + defer cleanupLayerTree(tree) - record, err := m.materializeLayerArtifact(desc) + record, err := m.materializeLayerArtifact(desc, tree) require.NoError(t, err) require.Contains(t, record.Whiteouts, whiteoutRecord{Dir: "gone", Target: "deleted.txt"}) diff --git a/lib/images/layer_gc.go b/lib/images/layer_gc.go new file mode 100644 index 000000000..aefdadfc0 --- /dev/null +++ b/lib/images/layer_gc.go @@ -0,0 +1,212 @@ +package images + +import ( + "context" + "fmt" + "io/fs" + "log/slog" + "os" + "path/filepath" + "strings" + "time" +) + +// layerEvictionGracePeriod keeps freshly written layer artifacts and temp +// directories out of cleanup so recovery and eviction never race builds that +// are still writing them. +const layerEvictionGracePeriod = 10 * time.Minute + +// referencedLayerDigests returns the set of layer blob digests referenced by +// the manifest models of every image in the content layout, plus the digests +// currently referenced by in-flight builds. Layer artifacts in this set are +// protected from eviction. Unreadable manifest models are skipped with a +// warning so one corrupt record cannot disable eviction entirely. +func (m *manager) referencedLayerDigests() map[string]struct{} { + refs := m.inflightLayerRefSnapshot() + contentRoot := filepath.Join(m.paths.ImagesDir(), "content") + err := filepath.WalkDir(contentRoot, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if entry.IsDir() || entry.Name() != "manifest.json" { + return nil + } + digestHex := filepath.Base(filepath.Dir(path)) + model, readErr := readManifestModel(m.paths, digestHex) + if readErr != nil { + slog.Warn("skipping unreadable manifest model for layer eviction", "digest", digestHex, "error", readErr) + return nil + } + if model == nil { + return nil + } + for _, layer := range model.Layers { + refs[strings.TrimPrefix(layer.Digest, "sha256:")] = struct{}{} + } + return nil + }) + if err != nil && !os.IsNotExist(err) { + slog.Warn("failed to walk content manifests for layer eviction", "error", err) + } + return refs +} + +// inflightLayerRefSnapshot returns the layer digests currently retained by +// in-flight builds. +func (m *manager) inflightLayerRefSnapshot() map[string]struct{} { + m.layerRefMu.Lock() + defer m.layerRefMu.Unlock() + refs := make(map[string]struct{}, len(m.inflightLayerRefs)) + for digestHex := range m.inflightLayerRefs { + refs[digestHex] = struct{}{} + } + return refs +} + +// reconcileLayerStore evicts unreferenced layer artifacts and refreshes the +// cached disk usage totals so accounting reflects the removals. +func (m *manager) reconcileLayerStore() { + m.evictUnreferencedLayerArtifacts() + m.refreshDiskUsageTotals() +} + +// evictUnreferencedLayerArtifacts removes layer artifacts that no image +// manifest model references, deleting the digest directory entirely. Artifacts +// newer than the grace period are kept so in-flight builds never lose work. +func (m *manager) evictUnreferencedLayerArtifacts() { + refs := m.referencedLayerDigests() + + layersDir := m.paths.ImageLayersDir() + entries, err := os.ReadDir(layersDir) + if err != nil { + if !os.IsNotExist(err) { + slog.Warn("layer eviction failed to list layer store", "error", err) + } + return + } + + cutoff := time.Now().Add(-m.layerEvictionGrace) + evicted := 0 + var evictedBytes int64 + for _, entry := range entries { + if !entry.IsDir() { + continue + } + digestHex := entry.Name() + if _, referenced := refs[digestHex]; referenced { + continue + } + size, removed := m.tryEvictLayerArtifact(digestHex, filepath.Join(layersDir, digestHex), cutoff) + if !removed { + continue + } + evicted++ + evictedBytes += size + } + if evicted > 0 { + slog.Info("evicted unreferenced layer artifacts", "count", evicted, "bytes", evictedBytes) + if m.metrics != nil { + m.metrics.layerArtifactsEvicted.Add(context.Background(), int64(evicted)) + } + } +} + +// tryEvictLayerArtifact removes one unreferenced layer artifact if it is still +// stale and no build is materializing it. The per-digest lock is taken with +// TryLock so eviction never blocks behind an in-flight conversion. +func (m *manager) tryEvictLayerArtifact(digestHex, dirPath string, cutoff time.Time) (int64, bool) { + lock := m.layerDigestLock(digestHex) + if !lock.TryLock() { + return 0, false + } + defer lock.Unlock() + + // The candidate was selected outside the lock; re-check that a build has + // not retained the digest and the artifact has not been rewritten since. + if _, referenced := m.inflightLayerRefSnapshot()[digestHex]; referenced { + return 0, false + } + info, statErr := os.Stat(dirPath) + if statErr != nil || info.ModTime().After(cutoff) { + return 0, false + } + size, err := dirSize(dirPath) + if err != nil { + slog.Warn("failed to measure layer artifact size", "digest", digestHex, "error", err) + } + if err := os.RemoveAll(dirPath); err != nil { + slog.Warn("failed to evict unreferenced layer artifact", "digest", digestHex, "error", err) + return 0, false + } + return size, true +} + +// cleanStaleImageTempDirs removes temp directories left behind by builds that +// were interrupted mid-install, mid-materialization, or mid-tag promotion. +// Only directories older than the grace period are removed so live builds are +// never disturbed. +func (m *manager) cleanStaleImageTempDirs() { + roots := []string{ + m.paths.ImageLayersDir(), + filepath.Join(m.paths.ImagesDir(), "content"), + } + cutoff := time.Now().Add(-m.layerEvictionGrace) + for _, root := range roots { + err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if !entry.IsDir() { + return nil + } + name := entry.Name() + if !strings.HasPrefix(name, ".unpack-") && !strings.HasPrefix(name, ".install-") && !strings.HasPrefix(name, ".tag-stage-") { + return nil + } + info, statErr := os.Stat(path) + if statErr == nil && info.ModTime().Before(cutoff) { + _ = os.RemoveAll(path) + } + return fs.SkipDir + }) + if err != nil && !os.IsNotExist(err) { + slog.Warn("failed to clean stale image temp dirs", "root", root, "error", err) + } + } +} + +// totalLayerArtifactBytes sums the bytes held by materialized layer +// artifacts, matching what diskutilization.Collect counts for the same store. +func totalLayerArtifactBytes(layersDir string) (int64, error) { + var total int64 + err := filepath.WalkDir(layersDir, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if entry.IsDir() { + return nil + } + if !strings.HasPrefix(entry.Name(), "layer.") { + return nil + } + info, statErr := entry.Info() + if statErr != nil { + return nil + } + total += info.Size() + return nil + }) + if err != nil && !os.IsNotExist(err) { + return 0, fmt.Errorf("walk layer artifacts: %w", err) + } + return total, nil +} diff --git a/lib/images/lifecycle_test.go b/lib/images/lifecycle_test.go new file mode 100644 index 000000000..1554f1501 --- /dev/null +++ b/lib/images/lifecycle_test.go @@ -0,0 +1,207 @@ +package images + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "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/kernel/hypeman/lib/paths" + "github.com/stretchr/testify/require" +) + +// writeSharedLayout writes several images into one OCI layout cache, each +// annotated with its own digest tag. +func writeSharedLayout(t *testing.T, p *paths.Paths, imgs ...gcr.Image) []string { + t.Helper() + + layoutPath, err := layout.Write(p.SystemOCICache(), empty.Index) + require.NoError(t, err) + + digests := make([]string, 0, len(imgs)) + for _, img := range imgs { + digest, err := img.Digest() + require.NoError(t, err) + require.NoError(t, layoutPath.AppendImage(img, layout.WithAnnotations(map[string]string{ + "org.opencontainers.image.ref.name": digestToLayoutTag(digest.String()), + }))) + digests = append(digests, digest.String()) + } + return digests +} + +func layerHexes(t *testing.T, p *paths.Paths) map[string]struct{} { + t.Helper() + entries, err := os.ReadDir(p.ImageLayersDir()) + require.NoError(t, err) + hexes := make(map[string]struct{}) + for _, entry := range entries { + if entry.IsDir() { + hexes[entry.Name()] = struct{}{} + } + } + return hexes +} + +// TestSharedLayersMaterializeOnceAndEvictWithReferences is the end-to-end +// lifecycle: two images share a base layer, the shared artifact is created +// once, survives the deletion of one image, and is evicted only when its last +// reference is gone. +func TestSharedLayersMaterializeOnceAndEvictWithReferences(t *testing.T) { + if _, err := exec.LookPath("mkfs.erofs"); err != nil { + t.Skip("mkfs.erofs not available") + } + dataDir := t.TempDir() + p := paths.New(dataDir) + mgr, err := NewManager(p, 1, nil) + require.NoError(t, err) + m := mgr.(*manager) + m.layerEvictionGrace = 0 + + base := syntheticLayer(t, "base.txt", "shared base content") + topA := syntheticLayer(t, "a.txt", "app A payload") + topB := syntheticLayer(t, "b.txt", "app B payload") + + imgA, err := mutate.AppendLayers(empty.Image, base, topA) + require.NoError(t, err) + imgB, err := mutate.AppendLayers(empty.Image, base, topB) + require.NoError(t, err) + + digests := writeSharedLayout(t, p, imgA, imgB) + digestA, digestB := digests[0], digests[1] + + baseManifest, err := imgA.Manifest() + require.NoError(t, err) + baseHex := baseManifest.Layers[0].Digest.Hex + topAHex := baseManifest.Layers[1].Digest.Hex + topBManifest, err := imgB.Manifest() + require.NoError(t, err) + topBHex := topBManifest.Layers[1].Digest.Hex + + ctx := context.Background() + const repoA = "kernel.local/apps/app-a" + const repoB = "kernel.local/apps/app-b" + + eventsA := make(chan StatusEvent, 2) + m.subscribeToReady(digestToLayoutTag(digestA), eventsA) + defer m.unsubscribeFromReady(digestToLayoutTag(digestA), eventsA) + _, err = m.ImportLocalImage(ctx, repoA, "v1", digestA) + require.NoError(t, err) + select { + case event := <-eventsA: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatal("image A did not become ready") + } + + eventsB := make(chan StatusEvent, 2) + m.subscribeToReady(digestToLayoutTag(digestB), eventsB) + defer m.unsubscribeFromReady(digestToLayoutTag(digestB), eventsB) + _, err = m.ImportLocalImage(ctx, repoB, "v1", digestB) + require.NoError(t, err) + select { + case event := <-eventsB: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatal("image B did not become ready") + } + + // The shared base layer materialized exactly once, alongside the two tops. + hexes := layerHexes(t, p) + require.Len(t, hexes, 3) + require.Contains(t, hexes, baseHex) + require.Contains(t, hexes, topAHex) + require.Contains(t, hexes, topBHex) + + // Deleting image A evicts only its unique layer; the shared base survives. + require.NoError(t, m.DeleteImage(ctx, repoA+"@"+digestA)) + hexes = layerHexes(t, p) + require.Len(t, hexes, 2) + require.Contains(t, hexes, baseHex, "shared base must survive while referenced") + require.Contains(t, hexes, topBHex) + require.NotContains(t, hexes, topAHex) + + // Deleting image B removes the last references: everything is evicted. + require.NoError(t, m.DeleteImage(ctx, repoB+"@"+digestB)) + hexes = layerHexes(t, p) + require.Empty(t, hexes, "unreferenced layer artifacts must be evicted") +} + +func TestTotalImageBytesIncludesLayerArtifacts(t *testing.T) { + p := paths.New(t.TempDir()) + m := newTestManager(p) + + digestHex := "cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01" + require.NoError(t, os.MkdirAll(p.ImageLayerDir(digestHex), 0o755)) + payload := make([]byte, 4096) + require.NoError(t, os.WriteFile(p.ImageLayerArtifact(digestHex), payload, 0o644)) + + totals, err := m.getDiskUsageTotals() + require.NoError(t, err) + require.GreaterOrEqual(t, totals.layerBytes, int64(len(payload))) + + totalBytes, err := m.TotalImageBytes(context.Background()) + require.NoError(t, err) + require.Equal(t, totals.readyImageBytes+totals.layerBytes, totalBytes) +} + +func TestCleanStaleImageTempDirsRemovesOnlyOldDirectories(t *testing.T) { + p := paths.New(t.TempDir()) + m := newTestManager(p) + m.layerEvictionGrace = time.Hour + + layersDir := p.ImageLayersDir() + staleDir := filepath.Join(layersDir, "ab12", ".unpack-stale") + freshDir := filepath.Join(layersDir, "cd34", ".unpack-fresh") + require.NoError(t, os.MkdirAll(staleDir, 0o755)) + require.NoError(t, os.MkdirAll(freshDir, 0o755)) + old := time.Now().Add(-2 * time.Hour) + require.NoError(t, os.Chtimes(staleDir, old, old)) + + m.cleanStaleImageTempDirs() + + _, err := os.Stat(staleDir) + require.True(t, os.IsNotExist(err), "stale temp dir must be removed") + _, err = os.Stat(freshDir) + require.NoError(t, err, "fresh temp dir must survive cleanup") +} + +func TestEvictionKeepsReferencedAndFreshArtifacts(t *testing.T) { + p := paths.New(t.TempDir()) + m := newTestManager(p) + m.layerEvictionGrace = time.Hour + + referencedHex := "ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01" + orphanFreshHex := "ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23" + + // A manifest model referencing one layer protects it regardless of age. + model := &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: "sha256:" + referencedHex, + Config: manifestConfigRef{ + Digest: "sha256:" + strings.Repeat("c", 64), + DiffIDs: []string{"sha256:" + referencedHex}, + }, + Layers: []layerDescriptor{{Digest: "sha256:" + referencedHex, DiffID: "sha256:" + referencedHex}}, + } + require.NoError(t, writeManifestModel(p, referencedHex, model)) + require.NoError(t, os.MkdirAll(p.ImageLayerDir(referencedHex), 0o755)) + require.NoError(t, os.WriteFile(p.ImageLayerArtifact(referencedHex), []byte("kept"), 0o644)) + + // An unreferenced but fresh artifact is protected by the grace period. + require.NoError(t, os.MkdirAll(p.ImageLayerDir(orphanFreshHex), 0o755)) + require.NoError(t, os.WriteFile(p.ImageLayerArtifact(orphanFreshHex), []byte("fresh"), 0o644)) + + m.reconcileLayerStore() + + hexes := layerHexes(t, p) + require.Contains(t, hexes, referencedHex) + require.Contains(t, hexes, orphanFreshHex) +} diff --git a/lib/images/manager.go b/lib/images/manager.go index a9b78c8b4..8b0fbb95b 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -76,13 +76,19 @@ type manager struct { queue *queue.Queue createMu sync.Mutex diskUsageMu sync.RWMutex + layerDigestMu sync.Mutex + layerDigestLocks map[string]*sync.Mutex tagGenerations map[string]uint64 + layerRefMu sync.Mutex + inflightLayerRefs map[string]int diskUsageLoaded bool readyImageBytes int64 + layerBytes int64 ociCacheBytes int64 metrics *Metrics inflightPulls map[string]*inflightImagePull // keyed by digest borrowedCredentialsTimeout time.Duration + layerEvictionGrace time.Duration readySubscribers map[string][]chan StatusEvent // keyed by digestHex subscriberMu sync.RWMutex } @@ -103,8 +109,11 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma queue: queue.New(maxConcurrentBuilds), inflightPulls: make(map[string]*inflightImagePull), borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, - readySubscribers: make(map[string][]chan StatusEvent), + layerEvictionGrace: layerEvictionGracePeriod, tagGenerations: make(map[string]uint64), + layerDigestLocks: make(map[string]*sync.Mutex), + inflightLayerRefs: make(map[string]int), + readySubscribers: make(map[string][]chan StatusEvent), } // Initialize metrics if meter is provided @@ -117,17 +126,24 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma } m.RecoverInterruptedBuilds() - m.promoteLegacyImages() + promoteLegacyImages(m.paths) + m.cleanStaleImageTempDirs() + m.reconcileLayerStore() return m, nil } -// promoteLegacyImages migrates ready legacy per-repository images into the -// shared content layout so every repository referencing the same digest -// shares one rootfs copy. Promotion hardlinks the disk, repoints tags, and -// removes the legacy tree; failures only warn so a partial migration never -// blocks startup. -func (m *manager) promoteLegacyImages() { - promoteLegacyImages(m.paths) +// layerDigestLock returns the mutex serializing materialization and eviction +// for one layer digest, so concurrent builds only contend on the digests they +// actually share. +func (m *manager) layerDigestLock(digestHex string) *sync.Mutex { + m.layerDigestMu.Lock() + defer m.layerDigestMu.Unlock() + lock := m.layerDigestLocks[digestHex] + if lock == nil { + lock = &sync.Mutex{} + m.layerDigestLocks[digestHex] = lock + } + return lock } func credentialsPresent(credentials *authn.AuthConfig) bool { @@ -257,8 +273,8 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest if meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()); err == nil { // Don't cache failed builds - allow retry 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) + if err := m.discardFailedImage(ref.Repository(), ref.DigestHex()); err != nil { + return nil, err } // Fall through to re-queue the build } else { @@ -293,8 +309,8 @@ func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthCo 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) + if err := m.discardFailedImage(ref.Repository(), ref.DigestHex()); err != nil { + return nil, true, err } return nil, false, nil } @@ -323,6 +339,16 @@ func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthCo return img, true, nil } +// discardFailedImage clears a failed build's tags and digest layout so a +// retry can start clean. +func (m *manager) discardFailedImage(repository, digestHex string) error { + _ = deleteTagsForDigest(m.paths, repository, digestHex) + if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, false); err != nil { + return fmt.Errorf("remove failed image: %w", err) + } + return nil +} + func tagGenerationKey(repository, tag string) string { return repository + ":" + tag } @@ -351,6 +377,13 @@ func (m *manager) restoreTagGenerations(metas []*imageMetadata) { } } +func (m *manager) claimReadyTag(ref *ResolvedRef) error { + m.createMu.Lock() + defer m.createMu.Unlock() + m.nextTagGeneration(ref.Repository(), ref.Tag()) + return createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) +} + func (m *manager) inflightCredentialsMatch(digest string, credentials *authn.AuthConfig) bool { inflight := m.inflightPulls[digest] var existingFingerprint [32]byte @@ -361,9 +394,6 @@ func (m *manager) inflightCredentialsMatch(digest string, credentials *authn.Aut } func (m *manager) registerInflightPull(digest string, credentials *authn.AuthConfig) *inflightImagePull { - if m.inflightPulls == nil { - m.inflightPulls = make(map[string]*inflightImagePull) - } if previous := m.inflightPulls[digest]; previous != nil && previous.timer != nil { previous.timer.Stop() } @@ -439,6 +469,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, return nil, fmt.Errorf("resolve existing image tag: %w", err) } } + tagGeneration := uint64(0) if ref.Tag() != "" { tagGeneration = m.nextTagGeneration(ref.Repository(), ref.Tag()) @@ -497,6 +528,31 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, return img, nil } +func (m *manager) retainLayerRefs(model *imageManifestModel) func() { + if model == nil { + return func() {} + } + m.layerRefMu.Lock() + digests := make([]string, 0, len(model.Layers)) + for _, layer := range model.Layers { + digestHex := strings.TrimPrefix(layer.Digest, "sha256:") + m.inflightLayerRefs[digestHex]++ + digests = append(digests, digestHex) + } + m.layerRefMu.Unlock() + + return func() { + m.layerRefMu.Lock() + defer m.layerRefMu.Unlock() + for _, digestHex := range digests { + m.inflightLayerRefs[digestHex]-- + if m.inflightLayerRefs[digestHex] == 0 { + delete(m.inflightLayerRefs, digestHex) + } + } + } +} + func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials *authn.AuthConfig, buildID string) { buildStart := time.Now() buildStatus := "failed" @@ -535,16 +591,22 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } m.recordPullMetrics(ctx, "success") + // The pulled layers own fully unpacked staging trees in /tmp; release them + // when the build ends no matter which path it takes. + defer result.cleanup() + + releaseLayerRefs := m.retainLayerRefs(result.Manifest) + defer func() { + releaseLayerRefs() + m.reconcileLayerStore() + }() + // Check if this digest already exists and is ready (deduplication) if meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()); err == nil { if meta.Status == StatusReady { // 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()) - m.createMu.Unlock() - if err != nil { + if err := m.claimReadyTag(ref); err != nil { slog.Warn("failed to claim ready image tag", "repository", ref.Repository(), "tag", ref.Tag(), "error", err) } } @@ -553,6 +615,9 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } } + // Materialization is best effort; the composed rootfs does not depend on it. + m.materializeLayerArtifacts(ctx, ref.Digest(), result) + m.updateStatusByDigest(ref, StatusConverting, nil, buildID) diskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk @@ -583,6 +648,27 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials buildStatus = "success" } +func (m *manager) materializeLayerArtifacts(ctx context.Context, digest string, result *pullResult) { + if result == nil || result.Manifest == nil { + return + } + start := time.Now() + var firstErr error + for _, desc := range result.Manifest.Layers { + if _, err := m.materializeLayerArtifact(desc, result.LayerTrees[desc.Digest]); err != nil { + slog.WarnContext(ctx, "failed to materialize layer artifact", "digest", desc.Digest, "error", err) + if firstErr == nil { + firstErr = err + } + } + } + cacheStatus := "miss" + if result.CacheHit { + cacheStatus = "hit" + } + m.recordImageBuildPhase(ctx, digest, "layer_materialization", time.Since(start), phaseStatus(firstErr), cacheStatus) +} + func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string) error { if diskTempPath != "" { defer os.Remove(diskTempPath) @@ -598,11 +684,6 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i } 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 @@ -617,12 +698,13 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i // 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) - } + if result.Manifest == nil { + return fmt.Errorf("manifest model missing for new image") + } + 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 @@ -635,24 +717,38 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i meta.Labels = result.Metadata.Labels meta.WorkingDir = result.Metadata.WorkingDir + if err := installAtomically(finalDiskPath, func(path string) error { + return os.Rename(diskTempPath, path) + }); err != nil { + _ = os.Remove(m.paths.ImageContentManifestModel(ref.DigestHex())) + return fmt.Errorf("install image disk: %w", err) + } if err := writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta); err != nil { + _ = os.Remove(finalDiskPath) + _ = os.Remove(m.paths.ImageContentManifestModel(ref.DigestHex())) return fmt.Errorf("write final metadata: %w", err) } m.notifyReady(ref.DigestHex(), StatusReady, nil) - if meta.RequestedTag != "" { - current, resolveErr := resolveTag(m.paths, ref.Repository(), meta.RequestedTag) - generationMatches := m.tagGenerations[tagGenerationKey(ref.Repository(), meta.RequestedTag)] == meta.TagGeneration - if resolveErr == nil && generationMatches && (current == ref.DigestHex() || current == meta.PreviousTagDigest) { - if err := createTagSymlink(m.paths, ref.Repository(), meta.RequestedTag, ref.DigestHex()); err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) - } - } - } + m.publishReadyTag(ref, meta) m.refreshDiskUsageTotals() return nil } +func (m *manager) publishReadyTag(ref *ResolvedRef, meta *imageMetadata) { + if meta.RequestedTag == "" { + return + } + currentDigest, err := resolveTag(m.paths, ref.Repository(), meta.RequestedTag) + generationMatches := m.tagGenerations[tagGenerationKey(ref.Repository(), meta.RequestedTag)] == meta.TagGeneration + if err != nil || !generationMatches || (currentDigest != ref.DigestHex() && currentDigest != meta.PreviousTagDigest) { + return + } + if err := createTagSymlink(m.paths, ref.Repository(), meta.RequestedTag, ref.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" @@ -708,7 +804,16 @@ func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err erro meta.Error = &errorMsg } - writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta) + if writeErr := writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta); writeErr != nil { + if status == StatusReady || status == StatusFailed { + m.notifyReady(ref.DigestHex(), status, errors.Join(err, writeErr)) + } + return + } + + if status == StatusFailed { + m.refreshDiskUsageTotals() + } // Notify while holding createMu so a delete/recreate cannot race the // metadata write and receive a terminal event for the old build. @@ -826,13 +931,22 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, 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("create image alias: %w", err) - } + + previousDigest := "" + if targetDigest, tagErr := resolveTag(m.paths, targetRef.Repository(), targetRef.Tag()); tagErr == nil { + previousDigest = targetDigest + } else if !errors.Is(tagErr, ErrNotFound) { + return nil, tagErr } - if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { - return nil, fmt.Errorf("create image tag: %w", err) + + if err := m.installImageTag(sourceRef, targetRef, digestHex, meta); err != nil { + return nil, err + } + if previousDigest != "" && previousDigest != digestHex { + if err := removeDigestIfUnreferenced(m.paths, targetRef.Repository(), previousDigest, true); err != nil { + return nil, fmt.Errorf("remove replaced image: %w", err) + } + m.reconcileLayerStore() } img := meta.toImage() @@ -840,6 +954,19 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, return img, nil } +func (m *manager) installImageTag(sourceRef, targetRef *NormalizedRef, digestHex string, meta *imageMetadata) error { + if sourceRef.Repository() != targetRef.Repository() { + if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta, &tagTarget{repository: targetRef.Repository(), tag: targetRef.Tag()}); err != nil { + return fmt.Errorf("create image alias: %w", err) + } + return nil + } + if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { + return fmt.Errorf("create image tag: %w", err) + } + return nil +} + func (m *manager) resolveTagSource(ref *NormalizedRef) (string, error) { if ref.IsDigest() { return ref.DigestHex(), nil @@ -883,7 +1010,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, false); err != nil { return err } - m.refreshDiskUsageTotals() + m.reconcileLayerStore() return nil } @@ -911,7 +1038,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, true); err != nil { return fmt.Errorf("delete orphaned digest %s: %w", digestHex, err) } - m.refreshDiskUsageTotals() + m.reconcileLayerStore() } return nil @@ -919,20 +1046,20 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { // TotalImageBytes returns the total size of all ready images on disk. func (m *manager) TotalImageBytes(ctx context.Context) (int64, error) { - readyImageBytes, _, err := m.getDiskUsageTotals() + totals, err := m.getDiskUsageTotals() if err != nil { return 0, err } - return readyImageBytes, nil + return totals.readyImageBytes + totals.layerBytes, nil } // TotalOCICacheBytes returns the total size of the OCI layer cache. func (m *manager) TotalOCICacheBytes(ctx context.Context) (int64, error) { - _, ociCacheBytes, err := m.getDiskUsageTotals() + totals, err := m.getDiskUsageTotals() if err != nil { return 0, err } - return ociCacheBytes, nil + return totals.ociCacheBytes, nil } func (m *manager) findRequestedTagImage(ref *NormalizedRef) *Image { @@ -959,6 +1086,10 @@ func (m *manager) findRequestedTagImage(ref *NormalizedRef) *Image { // 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 { ref, err := ParseNormalizedRef(name) if err != nil { diff --git a/lib/images/manager_test.go b/lib/images/manager_test.go index 8abd90a65..36015dd5b 100644 --- a/lib/images/manager_test.go +++ b/lib/images/manager_test.go @@ -7,6 +7,7 @@ import ( "os" "path/filepath" "strings" + "sync" "testing" "time" @@ -17,6 +18,19 @@ import ( "github.com/stretchr/testify/require" ) +// newTestManager returns a manager with the maps NewManager initializes, so +// tests can construct one directly without nil-map guards in the manager. +func newTestManager(p *paths.Paths) *manager { + return &manager{ + paths: p, + tagGenerations: make(map[string]uint64), + layerDigestLocks: make(map[string]*sync.Mutex), + inflightLayerRefs: make(map[string]int), + inflightPulls: make(map[string]*inflightImagePull), + readySubscribers: make(map[string][]chan StatusEvent), + } +} + func TestConversionFailedErr(t *testing.T) { t.Run("without detail", func(t *testing.T) { assert.EqualError(t, conversionFailedErr(nil, nil), "image conversion failed") diff --git a/lib/images/metrics.go b/lib/images/metrics.go index d860885b2..75164a6cd 100644 --- a/lib/images/metrics.go +++ b/lib/images/metrics.go @@ -11,11 +11,12 @@ import ( // Metrics holds the metrics instruments for image operations. type Metrics struct { - buildDuration metric.Float64Histogram - buildPhaseDuration metric.Float64Histogram - ociLayerCount metric.Int64Histogram - ociCompressedBytes metric.Int64Histogram - pullsTotal metric.Int64Counter + buildDuration metric.Float64Histogram + buildPhaseDuration metric.Float64Histogram + ociLayerCount metric.Int64Histogram + ociCompressedBytes metric.Int64Histogram + pullsTotal metric.Int64Counter + layerArtifactsEvicted metric.Int64Counter } // newMetrics creates and registers all image metrics. @@ -78,6 +79,14 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) { return nil, err } + layerArtifactsEvicted, err := meter.Int64Counter( + "hypeman_images_layer_artifacts_evicted_total", + metric.WithDescription("Total number of shared layer artifacts evicted after their last reference was removed"), + ) + if err != nil { + return nil, err + } + // Register observable gauges for queue length and total images buildQueueLength, err := meter.Int64ObservableGauge( "hypeman_images_build_queue_length", @@ -123,11 +132,12 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) { } return &Metrics{ - buildDuration: buildDuration, - buildPhaseDuration: buildPhaseDuration, - ociLayerCount: ociLayerCount, - ociCompressedBytes: ociCompressedBytes, - pullsTotal: pullsTotal, + buildDuration: buildDuration, + buildPhaseDuration: buildPhaseDuration, + ociLayerCount: ociLayerCount, + ociCompressedBytes: ociCompressedBytes, + pullsTotal: pullsTotal, + layerArtifactsEvicted: layerArtifactsEvicted, }, nil } diff --git a/lib/images/metrics_test.go b/lib/images/metrics_test.go index 93007b300..1f58ab48c 100644 --- a/lib/images/metrics_test.go +++ b/lib/images/metrics_test.go @@ -16,10 +16,8 @@ import ( func TestImageBuildPhaseMetrics(t *testing.T) { reader := otelmetric.NewManualReader() provider := otelmetric.NewMeterProvider(otelmetric.WithReader(reader)) - m := &manager{ - paths: paths.New(t.TempDir()), - queue: queue.New(1), - } + m := newTestManager(paths.New(t.TempDir())) + m.queue = queue.New(1) metrics, err := newMetrics(provider.Meter("test"), m) require.NoError(t, err) diff --git a/lib/images/oci.go b/lib/images/oci.go index 672dd9fd7..9517c24b0 100644 --- a/lib/images/oci.go +++ b/lib/images/oci.go @@ -195,6 +195,7 @@ func (c *ociClient) inspectDigestPlatformAuth(ctx context.Context, imageRef stri type pullResult struct { Metadata *containerMetadata Manifest *imageManifestModel + LayerTrees map[string]layerTree Digest string // sha256:abc123... CacheHit bool LayerCount int @@ -223,6 +224,11 @@ func (r *pullResult) measure(phase string, operation func() error) error { return err } +// cleanup removes the transient layer staging trees the result owns. +func (r *pullResult) cleanup() { + cleanupLayerTrees(r.LayerTrees) +} + func (c *ociClient) pullAndExport(ctx context.Context, imageRef, digest, exportDir string) (*pullResult, error) { return c.pullAndExportWithAuth(ctx, imageRef, digest, exportDir, nil) } @@ -286,7 +292,9 @@ func (c *ociClient) pullAndExportWithPlatformAuth(ctx context.Context, imageRef, if err := validateModelPairing(layoutTag, model); err != nil { return err } - return c.composeRootfs(exportDir, model.Layers) + var err error + result.LayerTrees, err = c.composeRootfsWithLayerTrees(exportDir, model.Layers) + return err } return c.unpackLayers(ctx, layoutTag, exportDir) }); err != nil { diff --git a/lib/images/oci_public.go b/lib/images/oci_public.go index 7d336745b..554e1967f 100644 --- a/lib/images/oci_public.go +++ b/lib/images/oci_public.go @@ -47,7 +47,10 @@ func (c *OCIClient) InspectManifestForLinux(ctx context.Context, imageRef string // PullAndUnpack pulls an OCI image and unpacks it to a directory (public for system manager). // Always targets Linux platform since hypeman VMs are Linux guests. func (c *OCIClient) PullAndUnpack(ctx context.Context, imageRef, digest, exportDir string) error { - _, err := c.client.pullAndExport(ctx, imageRef, digest, exportDir) + result, err := c.client.pullAndExport(ctx, imageRef, digest, exportDir) + if result != nil { + defer result.cleanup() + } if err != nil { return fmt.Errorf("pull and unpack: %w", err) } diff --git a/lib/images/recovery_regression_test.go b/lib/images/recovery_regression_test.go index b979939a9..c3d7577a4 100644 --- a/lib/images/recovery_regression_test.go +++ b/lib/images/recovery_regression_test.go @@ -43,12 +43,9 @@ func TestRecoverInterruptedBuildsCapturedFixtureMarksBuildFailed(t *testing.T) { client, err := newOCIClient(p.SystemOCICache()) require.NoError(t, err) - m := &manager{ - paths: p, - ociClient: client, - queue: queue.New(1), - readySubscribers: make(map[string][]chan StatusEvent), - } + m := newTestManager(p) + m.ociClient = client + m.queue = queue.New(1) m.RecoverInterruptedBuilds() @@ -65,7 +62,7 @@ func TestRecoverInterruptedBuildsCapturedFixtureMarksBuildFailed(t *testing.T) { require.NotNil(t, meta.Error) assert.Equal(t, recoveryFixtureDigest, meta.Digest) assert.Equal(t, StatusFailed, meta.Status) - assert.Contains(t, *meta.Error, "config rootfs.diff_ids has 0 entries but manifest has 1 layers") + assert.Contains(t, *meta.Error, "manifest model has 0 diff ids for 1 layers") } func copyRecoveryFixture(t *testing.T) string { diff --git a/lib/images/storage.go b/lib/images/storage.go index 5d4da61f1..d1598d821 100644 --- a/lib/images/storage.go +++ b/lib/images/storage.go @@ -233,7 +233,23 @@ func readMetadataAt(layout imageLayout) (*imageMetadata, error) { return &meta, nil } -func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata) error { +// tagTarget is a cross-repository tag to install while promoting an image +// into shared content storage. +type tagTarget struct { + repository string + tag string +} + +func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata, target *tagTarget) (retErr error) { + var installedTarget *stagedTagSymlink + defer func() { + if retErr == nil || installedTarget == nil { + return + } + if restoreErr := restoreSymlinkState(installedTarget.linkPath, installedTarget.previous); restoreErr != nil { + retErr = errors.Join(retErr, fmt.Errorf("restore target tag: %w", restoreErr)) + } + }() contentReady := false if contentMeta, err := readContentMetadata(p, digestHex); err == nil { contentReady = contentMeta.Status == StatusReady @@ -262,6 +278,19 @@ func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, s } } + // Install a cross-repository target before touching source references. This + // makes target installation failures leave the source layout unchanged. + if target != nil { + staged, err := stageTagSymlink(p, target.repository, target.tag, digestHex) + if err != nil { + return fmt.Errorf("stage target tag: %w", err) + } + installedTarget = &staged + if err := installStagedTag(p, &staged); err != nil { + return fmt.Errorf("install target tag: %w", err) + } + } + // A legacy source may still have tags pointing at its repository-local // digest directory. Move those references to the shared content before // removing the duplicate legacy tree. @@ -289,25 +318,16 @@ func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, s } staged = append(staged, ref) } - for i, ref := range staged { - if err := os.Rename(ref.tempPath, ref.linkPath); err != nil { - if rollbackErr := rollbackTagSymlinks(staged[:i]); rollbackErr != nil { + for i := range staged { + if err := installStagedTag(p, &staged[i]); err != nil { + if rollbackErr := rollbackTagSymlinks(staged[:i+1]); rollbackErr != nil { return errors.Join(fmt.Errorf("promote legacy tag: %w", err), rollbackErr) } return fmt.Errorf("promote legacy tag: %w", err) } } - staleClean := true - for _, ref := range staged { - if err := removeStaleTagSymlink(p, &ref); err != nil { - staleClean = false - fmt.Fprintf(os.Stderr, "Warning: failed to remove stale tag symlink %s: %v\n", ref.tag, err) - } - } - if staleClean { - if err := os.RemoveAll(legacyDir); err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to remove legacy digest directory %s: %v\n", digestHex, err) - } + if err := os.RemoveAll(legacyDir); err != nil { + fmt.Fprintf(os.Stderr, "Warning: failed to remove legacy digest directory %s: %v\n", digestHex, err) } } @@ -382,6 +402,19 @@ func stageTagSymlink(p *paths.Paths, repository, tag, digestHex string) (stagedT }, nil } +// installStagedTag atomically installs a staged tag symlink, removes the +// stale legacy-layout link it replaces, and cleans up the staging directory. +// A failure leaves the link's previous state recorded on ref available for +// restoreSymlinkState. +func installStagedTag(p *paths.Paths, ref *stagedTagSymlink) error { + err := os.Rename(ref.tempPath, ref.linkPath) + _ = os.RemoveAll(ref.tempDir) + if err != nil { + return err + } + return removeStaleTagSymlink(p, ref) +} + func readSymlinkState(path string) (symlinkState, error) { target, err := os.Readlink(path) if err != nil { @@ -443,14 +476,9 @@ func createTagSymlink(p *paths.Paths, repository, tag, digestHex string) error { if err != nil { return fmt.Errorf("stage tag symlink: %w", err) } - if err := os.Rename(ref.tempPath, ref.linkPath); err != nil { - _ = os.RemoveAll(ref.tempDir) + if err := installStagedTag(p, &ref); err != nil { 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) - } - _ = os.RemoveAll(ref.tempDir) return nil } diff --git a/lib/images/storage_refs.go b/lib/images/storage_refs.go index 5cba6b6f4..d901ba2af 100644 --- a/lib/images/storage_refs.go +++ b/lib/images/storage_refs.go @@ -10,6 +10,7 @@ import ( "github.com/kernel/hypeman/lib/paths" ) +// ensurePendingTag creates a pending tag only when no current tag exists. func ensurePendingTag(p *paths.Paths, repository, tag, digestHex string) error { _, err := resolveTag(p, repository, tag) if err == nil { @@ -89,7 +90,7 @@ func promoteLegacyImages(p *paths.Paths) { if readErr != nil || meta.Status != StatusReady { continue } - if promoteErr := promoteImageToContent(p, ref.repository, ref.digestHex, meta); promoteErr != nil { + if promoteErr := promoteImageToContent(p, ref.repository, ref.digestHex, meta, nil); promoteErr != nil { fmt.Fprintf(os.Stderr, "Warning: failed to promote legacy image %s@%s: %v\n", ref.repository, ref.digestHex, promoteErr) } } diff --git a/lib/images/storage_test.go b/lib/images/storage_test.go index 1bceb066f..cce34ec6c 100644 --- a/lib/images/storage_test.go +++ b/lib/images/storage_test.go @@ -247,8 +247,7 @@ func TestPromoteLegacyImagesMovesContentAndTags(t *testing.T) { require.NoError(t, os.MkdirAll(filepath.Dir(tagPath), 0o755)) require.NoError(t, os.Symlink(digest, tagPath)) - m := &manager{paths: p} - m.promoteLegacyImages() + promoteLegacyImages(p) // Content exists with the same bytes and is ready. contentMeta, err := readContentMetadata(p, digest) @@ -282,8 +281,7 @@ func TestPromoteLegacyImagesSkipsNonReady(t *testing.T) { CreatedAt: time.Now().UTC(), })) - m := &manager{paths: p} - m.promoteLegacyImages() + promoteLegacyImages(p) _, err := os.Stat(p.ImageContentMetadata(digest)) require.True(t, os.IsNotExist(err), "non-ready legacy image must not be promoted") diff --git a/lib/images/tag_test.go b/lib/images/tag_test.go index e44e32259..e082ad668 100644 --- a/lib/images/tag_test.go +++ b/lib/images/tag_test.go @@ -27,7 +27,7 @@ func seedReadyContentImage(t *testing.T, p *paths.Paths, repository, tag, digest } func newTagTestManager(p *paths.Paths) *manager { - return &manager{paths: p} + return newTestManager(p) } func TestTagImageSameRepository(t *testing.T) { @@ -180,9 +180,7 @@ func TestTagImageReplacesExistingTag(t *testing.T) { require.NoError(t, err) require.Equal(t, first, resolved) - // The orphaned second digest is cleaned up by delete semantics, not by tag; - // it is still referenced by no tag after the repoint only if it had no other - // tag. Here "stable" was its only tag, so it remains on disk until deleted. + // Replacing the only tag removes the old digest once it is unreferenced. _, err = os.Stat(p.ImageContentDir(second)) - require.NoError(t, err) + require.ErrorIs(t, err, os.ErrNotExist) }