From 8f956f2f8314eabe4a15a2f980646ef70d7af5c0 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 16:55:44 +0000 Subject: [PATCH 01/11] Add shared-layer lifecycle: accounting, eviction, recovery, and materialization Wire the shared-layer path into image builds: after a pull, each layer is materialized into the content-addressed artifact store (best effort), so layers shared across similar images are converted once and reused. A new layer_materialization build phase records the work. Reference-protected cleanup: deleting an image evicts only the layer artifacts no remaining manifest model references, with a grace period so in-flight builds never lose work. Startup recovery removes stale temp directories from interrupted materializations and installs, and evicts orphans left by an unclean shutdown. Layer-aware accounting: disk usage totals now include the physical bytes of the layer artifact store (TotalLayerBytes), and TotalImageBytes reports ready images plus shared layers so capacity admission sees the full footprint. Adds a hypeman_images_layer_artifacts_evicted_total counter. Legacy flattened images remain readable and bootable until retired; they simply have no manifest model and contribute no layer references. --- lib/images/disk_usage.go | 28 +++-- lib/images/layer_gc.go | 159 ++++++++++++++++++++++++++++ lib/images/lifecycle_test.go | 197 +++++++++++++++++++++++++++++++++++ lib/images/manager.go | 41 +++++++- lib/images/metrics.go | 30 ++++-- 5 files changed, 432 insertions(+), 23 deletions(-) create mode 100644 lib/images/layer_gc.go create mode 100644 lib/images/lifecycle_test.go diff --git a/lib/images/disk_usage.go b/lib/images/disk_usage.go index b65dba1ac..5a516d172 100644 --- a/lib/images/disk_usage.go +++ b/lib/images/disk_usage.go @@ -108,57 +108,65 @@ func totalOCICacheBlobBytesFromFilesystem(blobDir string) (int64, error) { return total, nil } -func (m *manager) getDiskUsageTotals() (int64, int64, error) { +func (m *manager) getDiskUsageTotals() (int64, int64, int64, error) { m.diskUsageMu.RLock() if m.diskUsageLoaded { readyImageBytes := m.readyImageBytes + layerBytes := m.layerBytes ociCacheBytes := m.ociCacheBytes m.diskUsageMu.RUnlock() - return readyImageBytes, ociCacheBytes, nil + return readyImageBytes, layerBytes, ociCacheBytes, nil } m.diskUsageMu.RUnlock() - readyImageBytes, ociCacheBytes, err := m.computeDiskUsageTotals() + readyImageBytes, layerBytes, ociCacheBytes, err := m.computeDiskUsageTotals() if err != nil { - return 0, 0, err + return 0, 0, 0, err } m.diskUsageMu.Lock() if !m.diskUsageLoaded { m.readyImageBytes = readyImageBytes + m.layerBytes = layerBytes m.ociCacheBytes = ociCacheBytes m.diskUsageLoaded = true } readyImageBytes = m.readyImageBytes + layerBytes = m.layerBytes ociCacheBytes = m.ociCacheBytes m.diskUsageMu.Unlock() - return readyImageBytes, ociCacheBytes, nil + return readyImageBytes, layerBytes, ociCacheBytes, nil } func (m *manager) refreshDiskUsageTotals() { - readyImageBytes, ociCacheBytes, err := m.computeDiskUsageTotals() + readyImageBytes, layerBytes, ociCacheBytes, err := m.computeDiskUsageTotals() if err != nil { return } m.diskUsageMu.Lock() m.readyImageBytes = readyImageBytes + m.layerBytes = layerBytes m.ociCacheBytes = ociCacheBytes m.diskUsageLoaded = true m.diskUsageMu.Unlock() } -func (m *manager) computeDiskUsageTotals() (int64, int64, error) { +func (m *manager) computeDiskUsageTotals() (int64, int64, int64, error) { readyImageBytes, err := totalReadyImageBytesFromMetadata(m.paths.ImagesDir()) if err != nil { - return 0, 0, err + return 0, 0, 0, err + } + layerBytes, err := totalLayerArtifactBytes(m.paths.ImageLayersDir()) + if err != nil { + return 0, 0, 0, err } ociCacheBytes, err := totalOCICacheBlobBytesFromFilesystem(m.paths.OCICacheBlobDir()) if err != nil { - return 0, 0, err + return 0, 0, 0, err } - return readyImageBytes, ociCacheBytes, nil + return readyImageBytes, layerBytes, ociCacheBytes, nil } func totalRootfsBytesInDigestDir(digestDir string) (int64, error) { diff --git a/lib/images/layer_gc.go b/lib/images/layer_gc.go new file mode 100644 index 000000000..9423d2f04 --- /dev/null +++ b/lib/images/layer_gc.go @@ -0,0 +1,159 @@ +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. Layer artifacts +// in this set are protected from eviction. +func (m *manager) referencedLayerDigests() (map[string]struct{}, error) { + refs := make(map[string]struct{}) + 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 || model == nil { + return nil + } + for _, layer := range model.Layers { + refs[strings.TrimPrefix(layer.Digest, "sha256:")] = struct{}{} + } + return nil + }) + if err != nil && !os.IsNotExist(err) { + return nil, fmt.Errorf("walk content manifests: %w", err) + } + return refs, nil +} + +// 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. +// Returns the number of artifacts and bytes removed. +func (m *manager) evictUnreferencedLayerArtifacts() (int, int64) { + refs, err := m.referencedLayerDigests() + if err != nil { + slog.Warn("layer eviction skipped", "error", err) + return 0, 0 + } + + 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 0, 0 + } + + 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 + } + dirPath := filepath.Join(layersDir, digestHex) + info, statErr := os.Stat(dirPath) + if statErr != nil || info.ModTime().After(cutoff) { + continue + } + size, _ := dirSize(dirPath) + if removeErr := os.RemoveAll(dirPath); removeErr != nil { + slog.Warn("failed to evict unreferenced layer artifact", "digest", digestHex, "error", removeErr) + 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)) + } + } + return evicted, evictedBytes +} + +// cleanStaleImageTempDirs removes temp directories left behind by builds that +// were interrupted mid-install or mid-materialization. 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 { + filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + return nil + } + if !entry.IsDir() { + return nil + } + name := entry.Name() + if !strings.HasPrefix(name, ".unpack-") && !strings.HasPrefix(name, ".install-") { + return nil + } + info, statErr := os.Stat(path) + if statErr == nil && info.ModTime().Before(cutoff) { + _ = os.RemoveAll(path) + } + return fs.SkipDir + }) + } +} + +// totalLayerArtifactBytes sums the physical bytes held by the shared layer +// artifact 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 + } + 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..dda2f1e0f --- /dev/null +++ b/lib/images/lifecycle_test.go @@ -0,0 +1,197 @@ +package images + +import ( + "context" + "os" + "path/filepath" + "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) { + 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 := &manager{paths: 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)) + + layerBytes, err := m.TotalLayerBytes(context.Background()) + require.NoError(t, err) + require.GreaterOrEqual(t, layerBytes, int64(len(payload))) + + totalBytes, err := m.TotalImageBytes(context.Background()) + require.NoError(t, err) + require.Equal(t, layerBytes, totalBytes, "no ready images, total must equal layer bytes") +} + +func TestCleanStaleImageTempDirsRemovesOnlyOldDirectories(t *testing.T) { + p := paths.New(t.TempDir()) + m := &manager{paths: p, 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 := &manager{paths: p, 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, + Layers: []layerDescriptor{{Digest: "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)) + + evicted, _ := m.evictUnreferencedLayerArtifacts() + require.Equal(t, 0, evicted) + + 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..681d1e518 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -79,10 +79,12 @@ type manager struct { tagGenerations map[string]uint64 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,6 +105,7 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma queue: queue.New(maxConcurrentBuilds), inflightPulls: make(map[string]*inflightImagePull), borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, + layerEvictionGrace: layerEvictionGracePeriod, readySubscribers: make(map[string][]chan StatusEvent), tagGenerations: make(map[string]uint64), } @@ -118,6 +121,9 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma m.RecoverInterruptedBuilds() m.promoteLegacyImages() + m.cleanStaleImageTempDirs() + m.evictUnreferencedLayerArtifacts() + m.refreshDiskUsageTotals() return m, nil } @@ -553,6 +559,23 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } } + // Materialize shared per-layer artifacts for the pulled image. Best effort: + // the composed rootfs is built from blobs, so an artifact failure degrades + // layer sharing and accounting but does not make the image unusable. + if result.Manifest != nil { + materializeStart := time.Now() + for _, desc := range result.Manifest.Layers { + if _, err := m.materializeLayerArtifact(desc); err != nil { + slog.WarnContext(ctx, "failed to materialize layer artifact", "digest", desc.Digest, "error", err) + } + } + cacheStatus := "miss" + if result.CacheHit { + cacheStatus = "hit" + } + m.recordImageBuildPhase(ctx, ref.Digest(), "layer_materialization", time.Since(materializeStart), "success", cacheStatus) + } + m.updateStatusByDigest(ref, StatusConverting, nil, buildID) diskPath := resolveImageLayout(m.paths, ref.Repository(), ref.DigestHex()).disk @@ -883,6 +906,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, false); err != nil { return err } + m.evictUnreferencedLayerArtifacts() m.refreshDiskUsageTotals() return nil } @@ -911,6 +935,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.evictUnreferencedLayerArtifacts() m.refreshDiskUsageTotals() } @@ -919,16 +944,26 @@ 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() + readyImageBytes, layerBytes, _, err := m.getDiskUsageTotals() + if err != nil { + return 0, err + } + return readyImageBytes + layerBytes, nil +} + +// TotalLayerBytes returns the physical bytes held by the shared layer +// artifact store. +func (m *manager) TotalLayerBytes(ctx context.Context) (int64, error) { + _, layerBytes, _, err := m.getDiskUsageTotals() if err != nil { return 0, err } - return readyImageBytes, nil + return 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() + _, _, ociCacheBytes, err := m.getDiskUsageTotals() if err != nil { return 0, err } 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 } From 1786c7a215caadb7a669d42e513f33a2d088d282 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:20:25 +0000 Subject: [PATCH 02/11] Harden shared image storage lifecycle --- lib/diskutilization/diskutilization.go | 28 +++- lib/diskutilization/diskutilization_test.go | 16 +++ lib/images/compose.go | 83 +++++++++--- lib/images/layer_artifact.go | 137 ++++++++------------ lib/images/layer_gc.go | 13 +- lib/images/lifecycle_test.go | 7 +- lib/images/oci.go | 5 +- lib/images/storage.go | 19 ++- lib/images/storage_refs.go | 1 + lib/images/tag_test.go | 6 +- 10 files changed, 198 insertions(+), 117 deletions(-) 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..352280751 100644 --- a/lib/images/compose.go +++ b/lib/images/compose.go @@ -7,8 +7,9 @@ import ( "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 +18,91 @@ 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 } +// applyLayerToDir extracts one layer into a private staging directory, then +// applies whiteouts before copying the layer's entries onto the composed rootfs. func (c *ociClient) applyLayerToDir(dest string, desc layerDescriptor) error { + tree, err := c.extractLayerTree(desc) + if err != nil { + return err + } + defer cleanupLayerTree(tree) + if err := applyLayerTree(tree.path, dest); err != nil { + return fmt.Errorf("apply layer tree: %w", err) + } + return nil +} + +func (c *ociClient) extractLayerTree(desc layerDescriptor) (layerTree, 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) + return layerTree{}, fmt.Errorf("invalid layer digest: %s", desc.Digest) } 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) + cleanupLayerTree(layerTree{path: layerDir}) + return layerTree{}, 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) - } - return nil + return layerTree{path: layerDir, stats: stats}, nil } diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index 41d7201e8..83ef0a88c 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -63,10 +63,7 @@ type whiteoutRecord struct { } func (a *layerArtifact) matches(desc layerDescriptor) bool { - if a.Digest != desc.Digest || a.Format != layerArtifactFormat() { - return false - } - return desc.DiffID == "" || a.DiffID == desc.DiffID + return a.Digest == desc.Digest && a.Format == layerArtifactFormat() } const ( @@ -104,9 +101,6 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { if err := json.Unmarshal(data, &record); err != nil { return nil, fmt.Errorf("unmarshal layer record: %w", err) } - if err := record.validate(); err != nil { - return nil, fmt.Errorf("invalid layer record: %w", err) - } return &record, nil } @@ -115,31 +109,15 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { // 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 { - return fmt.Errorf("invalid digest or format") - } - if a.SizeBytes < 0 || a.UnpackedBytes < 0 || a.Entries < 0 { - return fmt.Errorf("invalid size or entry counts") - } - if a.Format == layerFormatExt4 && a.Options.Compression != "" { - return fmt.Errorf("ext4 artifact has compression options") - } - if a.Format == layerFormatErofs && a.Options.Compression != "lz4" { - return fmt.Errorf("erofs artifact has invalid compression options") - } - 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) } + m.layerLifecycleMu.Lock() + defer m.layerLifecycleMu.Unlock() + if record, err := readLayerRecord(m.paths, layerHex); err != nil { return nil, err } else if record != nil && record.matches(desc) { @@ -178,11 +156,28 @@ func (m *manager) materializeLayerArtifact(desc layerDescriptor) (*layerArtifact return m.installLayerArtifact(desc, layerHex, unpackDir, stats) } -func artifactOptions(format string) layerArtifactOptions { - if format == layerFormatErofs { - return layerArtifactOptions{Compression: "lz4"} +func (m *manager) materializeLayerArtifactFromTree(desc layerDescriptor, tree layerTree) (*layerArtifact, error) { + layerHex := strings.TrimPrefix(desc.Digest, "sha256:") + if layerHex == "" || strings.Contains(layerHex, "/") || layerHex == "." || strings.Contains(layerHex, "..") { + return nil, fmt.Errorf("invalid layer digest: %s", desc.Digest) + } + if tree.path == "" || tree.stats == nil { + return nil, fmt.Errorf("missing layer staging tree: %s", desc.Digest) + } + 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) + } + + m.layerLifecycleMu.Lock() + defer m.layerLifecycleMu.Unlock() + if record, err := readLayerRecord(m.paths, layerHex); 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 + } } - return layerArtifactOptions{} + return m.installLayerArtifact(desc, layerHex, tree.path, tree.stats) } func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir string, stats *unpackStats) (*layerArtifact, error) { @@ -191,7 +186,7 @@ func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir Digest: desc.Digest, DiffID: desc.DiffID, Format: layerArtifactFormat(), - Options: artifactOptions(layerArtifactFormat()), + Options: layerArtifactOptions{Compression: "lz4"}, UnpackedBytes: stats.unpackedBytes, Entries: stats.entries, Whiteouts: stats.whiteouts, @@ -227,6 +222,11 @@ func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir return record, nil } +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) { @@ -562,10 +574,7 @@ func applyLayerTree(layerDir, targetDir string) error { if err != nil { return err } - targetParent, err := safeJoin(targetDir, filepath.Dir(rel)) - if err != nil { - return err - } + targetParent := filepath.Join(targetDir, filepath.Dir(rel)) if base == opaqueWhiteout { return clearDirContents(targetParent) } @@ -625,7 +634,7 @@ func clearDirContents(dir string) error { return err } if !info.IsDir() { - return removePath(dir) + return fmt.Errorf("opaque whiteout target is not a directory: %s", dir) } entries, err := os.ReadDir(dir) if err != nil { @@ -659,7 +668,7 @@ func copyEntryInto(src, dst string, hardlinks map[hardlinkIdentity]string) error case 0: return copyRegularEntry(src, dst, info, hardlinks) case fs.ModeDir: - return copyDirectoryEntry(src, dst, info) + return copyDirectoryEntry(dst, info) case fs.ModeSymlink: return copySymlinkEntry(src, dst) default: @@ -681,10 +690,10 @@ func copyRegularEntry(src, dst string, info os.FileInfo, hardlinks map[hardlinkI if err := copyFileContents(src, dst); err != nil { return err } - return copyEntryMetadata(src, dst, info) + return copyEntryMetadata(dst, info) } -func copyDirectoryEntry(src, dst string, info os.FileInfo) error { +func copyDirectoryEntry(dst string, info os.FileInfo) error { if existing, err := os.Lstat(dst); err == nil && !existing.IsDir() { if err := removePath(dst); err != nil { return err @@ -693,7 +702,7 @@ func copyDirectoryEntry(src, dst string, info os.FileInfo) error { if err := os.MkdirAll(dst, info.Mode().Perm()); err != nil { return err } - return copyEntryMetadata(src, dst, info) + return copyEntryMetadata(dst, info) } func copySymlinkEntry(src, dst string) error { @@ -725,7 +734,7 @@ func copySpecialEntry(src, dst string, info os.FileInfo) error { if err := unix.Mknod(dst, mode|uint32(info.Mode().Perm()), int(stat.Rdev)); err != nil { return err } - return copyEntryMetadata(src, dst, info) + return copyEntryMetadata(dst, info) } func specialFileMode(mode fs.FileMode) (uint32, error) { @@ -741,61 +750,19 @@ func specialFileMode(mode fs.FileMode) (uint32, error) { } } -func copyEntryMetadata(src, dst string, info os.FileInfo) error { +func copyEntryMetadata(dst string, info os.FileInfo) error { if stat, ok := info.Sys().(*syscall.Stat_t); ok { if err := os.Lchown(dst, int(stat.Uid), int(stat.Gid)); err != nil && !errors.Is(err, os.ErrPermission) && !errors.Is(err, unix.EPERM) { return err } } if info.Mode()&os.ModeSymlink == 0 { - mode := info.Mode().Perm() | info.Mode()&(os.ModeSetuid|os.ModeSetgid|os.ModeSticky) - if err := os.Chmod(dst, mode); err != nil { + if err := os.Chmod(dst, info.Mode().Perm()); err != nil { return err } if err := os.Chtimes(dst, info.ModTime(), info.ModTime()); err != nil { return err } - if err := copyXattrs(src, dst); err != nil { - return err - } - } - return nil -} - -func copyXattrs(src, dst string) error { - size, err := unix.Llistxattr(src, nil) - if err != nil { - if errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EPERM) { - return nil - } - return err - } - names := make([]byte, size) - if size > 0 { - n, err := unix.Llistxattr(src, names) - if err != nil { - return err - } - names = names[:n] - } - for _, name := range strings.Split(strings.TrimSuffix(string(names), "\\x00"), "\\x00") { - if name == "" { - continue - } - size, err := unix.Lgetxattr(src, name, nil) - if err != nil { - if errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EPERM) || errors.Is(err, unix.ENODATA) { - continue - } - return err - } - value := make([]byte, size) - if _, err := unix.Lgetxattr(src, name, value); err != nil { - return err - } - if err := unix.Lsetxattr(dst, name, value, 0); err != nil && !errors.Is(err, unix.ENOTSUP) && !errors.Is(err, unix.EPERM) { - return err - } } return nil } diff --git a/lib/images/layer_gc.go b/lib/images/layer_gc.go index 9423d2f04..2cb25b880 100644 --- a/lib/images/layer_gc.go +++ b/lib/images/layer_gc.go @@ -34,7 +34,10 @@ func (m *manager) referencedLayerDigests() (map[string]struct{}, error) { } digestHex := filepath.Base(filepath.Dir(path)) model, readErr := readManifestModel(m.paths, digestHex) - if readErr != nil || model == nil { + if readErr != nil { + return fmt.Errorf("read manifest model %s: %w", digestHex, readErr) + } + if model == nil { return nil } for _, layer := range model.Layers { @@ -45,6 +48,11 @@ func (m *manager) referencedLayerDigests() (map[string]struct{}, error) { if err != nil && !os.IsNotExist(err) { return nil, fmt.Errorf("walk content manifests: %w", err) } + m.layerRefMu.Lock() + for digestHex := range m.inflightLayerRefs { + refs[digestHex] = struct{}{} + } + m.layerRefMu.Unlock() return refs, nil } @@ -53,6 +61,9 @@ func (m *manager) referencedLayerDigests() (map[string]struct{}, error) { // newer than the grace period are kept so in-flight builds never lose work. // Returns the number of artifacts and bytes removed. func (m *manager) evictUnreferencedLayerArtifacts() (int, int64) { + m.layerLifecycleMu.Lock() + defer m.layerLifecycleMu.Unlock() + refs, err := m.referencedLayerDigests() if err != nil { slog.Warn("layer eviction skipped", "error", err) diff --git a/lib/images/lifecycle_test.go b/lib/images/lifecycle_test.go index dda2f1e0f..23f09f2ee 100644 --- a/lib/images/lifecycle_test.go +++ b/lib/images/lifecycle_test.go @@ -4,6 +4,7 @@ import ( "context" "os" "path/filepath" + "strings" "testing" "time" @@ -178,7 +179,11 @@ func TestEvictionKeepsReferencedAndFreshArtifacts(t *testing.T) { model := &imageManifestModel{ SchemaVersion: manifestModelSchemaVersion, Digest: "sha256:" + referencedHex, - Layers: []layerDescriptor{{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)) diff --git a/lib/images/oci.go b/lib/images/oci.go index 672dd9fd7..29df90c04 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 @@ -286,7 +287,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/storage.go b/lib/images/storage.go index 5d4da61f1..7e060d9ed 100644 --- a/lib/images/storage.go +++ b/lib/images/storage.go @@ -233,7 +233,7 @@ func readMetadataAt(layout imageLayout) (*imageMetadata, error) { return &meta, nil } -func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata) error { +func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata, target ...string) error { contentReady := false if contentMeta, err := readContentMetadata(p, digestHex); err == nil { contentReady = contentMeta.Status == StatusReady @@ -262,6 +262,23 @@ 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 len(target) == 2 { + staged, err := stageTagSymlink(p, target[0], target[1], digestHex) + if err != nil { + return fmt.Errorf("stage target tag: %w", err) + } + if err := os.Rename(staged.tempPath, staged.linkPath); err != nil { + _ = os.RemoveAll(staged.tempDir) + return fmt.Errorf("install target tag: %w", err) + } + if err := removeStaleTagSymlink(p, &staged); err != nil { + return fmt.Errorf("remove stale target tag: %w", err) + } + _ = os.RemoveAll(staged.tempDir) + } + // 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. diff --git a/lib/images/storage_refs.go b/lib/images/storage_refs.go index 5cba6b6f4..e0ff05334 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 { diff --git a/lib/images/tag_test.go b/lib/images/tag_test.go index e44e32259..818ef8db3 100644 --- a/lib/images/tag_test.go +++ b/lib/images/tag_test.go @@ -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) } From 26dbf745228ce7f5b491262b832ad90fa72a1a25 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:46:53 +0000 Subject: [PATCH 03/11] Simplify layer materialization flow --- lib/images/manager.go | 374 ++++++++++++++++++++++-------------------- 1 file changed, 192 insertions(+), 182 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index 681d1e518..9d67d4889 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -76,7 +76,9 @@ type manager struct { queue *queue.Queue createMu sync.Mutex diskUsageMu sync.RWMutex - tagGenerations map[string]uint64 + layerLifecycleMu sync.Mutex + layerRefMu sync.Mutex + inflightLayerRefs map[string]int diskUsageLoaded bool readyImageBytes int64 layerBytes int64 @@ -106,8 +108,8 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma inflightPulls: make(map[string]*inflightImagePull), borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, layerEvictionGrace: layerEvictionGracePeriod, + inflightLayerRefs: make(map[string]int), readySubscribers: make(map[string][]chan StatusEvent), - tagGenerations: make(map[string]uint64), } // Initialize metrics if meter is provided @@ -120,18 +122,13 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma } m.RecoverInterruptedBuilds() - m.promoteLegacyImages() + promoteLegacyImages(m.paths) m.cleanStaleImageTempDirs() m.evictUnreferencedLayerArtifacts() m.refreshDiskUsageTotals() 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) } @@ -227,9 +224,49 @@ 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); 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 + if meta.Status == StatusFailed { + // Remove a failed tag before retrying. A tag for another digest is + // intentionally left untouched until the replacement is ready. + _ = deleteTagsForDigest(m.paths, ref.Repository(), ref.DigestHex()) + // Clean up the failed build directory so we can retry. Shared content + // is retained if another tag or pull still references it. + if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { + return nil, fmt.Errorf("remove failed image: %w", err) + } + // Fall through to re-queue the build + } else { + // Keep an existing ready tag visible until a replacement is ready. A + // pending tag is only created when it does not already point elsewhere. + if ref.Tag() != "" { + var tagErr error + if meta.Status == StatusReady { + tagErr = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } else { + tagErr = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } + if tagErr != nil { + return nil, fmt.Errorf("create image tag: %w", tagErr) + } + } + img := meta.toImage() + img.Name = ref.String() + if meta.Status == StatusReady { + return img, nil + } + if !m.inflightCredentialsMatch(ref.Digest(), req.Credentials) { + return nil, fmt.Errorf("%w: retry after the current pull completes", ErrCredentialConflict) + } + 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, req, platform) } @@ -263,6 +300,7 @@ 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 { + _ = deleteTagsForDigest(m.paths, ref.Repository(), ref.DigestHex()) if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { return nil, fmt.Errorf("remove failed image: %w", err) } @@ -271,7 +309,6 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest if ref.Tag() != "" { var tagErr error if meta.Status == StatusReady { - m.nextTagGeneration(ref.Repository(), ref.Tag()) tagErr = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) } else { tagErr = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) @@ -293,70 +330,6 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest return m.createAndQueueImage(ref, CreateImageRequest{Name: imageRef}, hostPlatform()) } -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 - } - if meta.Status == StatusFailed { - if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { - return nil, true, fmt.Errorf("remove failed image: %w", err) - } - return nil, false, nil - } - if ref.Tag() != "" { - if meta.Status == StatusReady { - m.nextTagGeneration(ref.Repository(), ref.Tag()) - err = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } else { - err = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } - if err != nil { - return nil, true, fmt.Errorf("create image tag: %w", err) - } - } - img := meta.toImage() - img.Name = 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 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) restoreTagGenerations(metas []*imageMetadata) { - m.createMu.Lock() - defer m.createMu.Unlock() - for _, meta := range metas { - if meta.RequestedTag == "" { - 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 - } - } -} - func (m *manager) inflightCredentialsMatch(digest string, credentials *authn.AuthConfig) bool { inflight := m.inflightPulls[digest] var existingFingerprint [32]byte @@ -445,10 +418,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()) - } + meta := &imageMetadata{ Name: ref.String(), Digest: ref.Digest(), @@ -460,7 +430,6 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, Tags: tags.Clone(req.Tags), RequestedTag: ref.Tag(), PreviousTagDigest: previousTagDigest, - TagGeneration: tagGeneration, CreatedAt: time.Now(), } @@ -503,6 +472,34 @@ 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() + if m.inflightLayerRefs == nil { + m.inflightLayerRefs = make(map[string]int) + } + 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" @@ -541,40 +538,27 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } m.recordPullMetrics(ctx, "success") + releaseLayerRefs := m.retainLayerRefs(result.Manifest) + defer func() { + releaseLayerRefs() + m.evictUnreferencedLayerArtifacts() + m.refreshDiskUsageTotals() + }() + // 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 { - slog.Warn("failed to claim ready image tag", "repository", ref.Repository(), "tag", ref.Tag(), "error", err) - } + createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) } buildStatus = "success" return } } - // Materialize shared per-layer artifacts for the pulled image. Best effort: - // the composed rootfs is built from blobs, so an artifact failure degrades - // layer sharing and accounting but does not make the image unusable. - if result.Manifest != nil { - materializeStart := time.Now() - for _, desc := range result.Manifest.Layers { - if _, err := m.materializeLayerArtifact(desc); err != nil { - slog.WarnContext(ctx, "failed to materialize layer artifact", "digest", desc.Digest, "error", err) - } - } - cacheStatus := "miss" - if result.CacheHit { - cacheStatus = "hit" - } - m.recordImageBuildPhase(ctx, ref.Digest(), "layer_materialization", time.Since(materializeStart), "success", cacheStatus) - } + // Materialization is best effort; the composed rootfs does not depend on it. + m.materializeLayerArtifacts(ctx, ref.Digest(), result) m.updateStatusByDigest(ref, StatusConverting, nil, buildID) @@ -606,6 +590,30 @@ 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 + } + defer cleanupLayerTrees(result.LayerTrees) + start := time.Now() + for _, desc := range result.Manifest.Layers { + var err error + if tree, ok := result.LayerTrees[desc.Digest]; ok { + _, err = m.materializeLayerArtifactFromTree(desc, tree) + } else { + _, err = m.materializeLayerArtifact(desc) + } + if err != nil { + slog.WarnContext(ctx, "failed to materialize layer artifact", "digest", desc.Digest, "error", err) + } + } + cacheStatus := "miss" + if result.CacheHit { + cacheStatus = "hit" + } + m.recordImageBuildPhase(ctx, digest, "layer_materialization", time.Since(start), "success", cacheStatus) +} + func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string) error { if diskTempPath != "" { defer os.Remove(diskTempPath) @@ -663,11 +671,11 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i } 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 { + if ref.Tag() != "" { + currentDigest, tagErr := resolveTag(m.paths, ref.Repository(), ref.Tag()) + shouldPublish := tagErr == nil && (currentDigest == ref.DigestHex() || currentDigest == meta.PreviousTagDigest) + if shouldPublish { + if err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()); err != nil { fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) } } @@ -731,7 +739,13 @@ 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 { + 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. @@ -745,7 +759,6 @@ 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 { @@ -839,48 +852,48 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, m.createMu.Lock() defer m.createMu.Unlock() - m.nextTagGeneration(targetRef.Repository(), targetRef.Tag()) - digestHex, err := m.resolveTagSource(sourceRef) - if err != nil { - return nil, err + digestHex := sourceRef.DigestHex() + if !sourceRef.IsDigest() { + digestHex, err = resolveTag(m.paths, sourceRef.Repository(), sourceRef.Tag()) + if err != nil { + return nil, err + } } - meta, err := m.readyImageMetadata(sourceRef.Repository(), digestHex) + + meta, err := readMetadata(m.paths, sourceRef.Repository(), digestHex) if err != nil { return nil, err } + if meta.Status != StatusReady { + return nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) + } + + 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 sourceRef.Repository() != targetRef.Repository() { - if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta); err != nil { + if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta, targetRef.Repository(), targetRef.Tag()); err != nil { return nil, fmt.Errorf("create image alias: %w", err) } - } - if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { + } else if err := createTagSymlink(m.paths, targetRef.Repository(), targetRef.Tag(), digestHex); err != nil { return nil, fmt.Errorf("create image tag: %w", 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) + } + } img := meta.toImage() img.Name = targetRef.String() return img, nil } -func (m *manager) resolveTagSource(ref *NormalizedRef) (string, error) { - if ref.IsDigest() { - return ref.DigestHex(), nil - } - return resolveTag(m.paths, ref.Repository(), ref.Tag()) -} - -func (m *manager) readyImageMetadata(repository, digestHex string) (*imageMetadata, error) { - meta, err := readMetadata(m.paths, repository, digestHex) - if err != nil { - return nil, err - } - if meta.Status != StatusReady { - return nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) - } - return meta, nil -} - func (m *manager) DeleteImage(ctx context.Context, name string) error { // Parse and normalize the reference ref, err := ParseNormalizedRef(name) @@ -912,7 +925,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) @@ -994,17 +1006,50 @@ 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 { return fmt.Errorf("parse image name: %w", err) } - img, err := m.waitForImage(ctx, name, ref) - if err != nil { - return err + + // Wait for the image to appear in the store. In the build flow, the + // registry triggers ImportLocalImage asynchronously after a push, so the + // image may not exist when the build manager calls WaitForReady. + const maxWaitForExist = 30 * time.Second + const pollInterval = 100 * time.Millisecond + + var img *Image + deadline := time.Now().Add(maxWaitForExist) + for { + if !ref.IsDigest() { + img = m.findRequestedTagImage(ref) + } + if img == nil { + img, err = m.GetImage(ctx, name) + } + if img != nil { + break + } + if time.Now().After(deadline) { + return fmt.Errorf("get image: %w", err) + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(pollInterval): + } } - if terminal, err := terminalImageError(img); terminal { - return err + + // Check if already in terminal state + switch img.Status { + case StatusReady: + return nil + case StatusFailed: + return conversionFailedErr(img.Error, nil) } digestHex := strings.TrimPrefix(img.Digest, "sha256:") @@ -1017,8 +1062,11 @@ 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 terminal, terminalErr := terminalImageError(img); terminal { - return terminalErr + switch img.Status { + case StatusReady: + return nil + case StatusFailed: + return conversionFailedErr(img.Error, nil) } } @@ -1034,44 +1082,6 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { } } -func (m *manager) waitForImage(ctx context.Context, name string, ref *NormalizedRef) (*Image, error) { - const maxWaitForExist = 30 * time.Second - const pollInterval = 100 * time.Millisecond - var lastErr error - deadline := time.Now().Add(maxWaitForExist) - for { - var img *Image - if !ref.IsDigest() { - img = m.findRequestedTagImage(ref) - } - if img == nil { - img, lastErr = m.GetImage(ctx, name) - } - if img != nil { - return img, nil - } - if time.Now().After(deadline) { - return nil, fmt.Errorf("get image: %w", lastErr) - } - select { - case <-ctx.Done(): - return nil, ctx.Err() - case <-time.After(pollInterval): - } - } -} - -func 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) From 16c40eef3a028ae4a0207f571375fcd5cb604461 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:57:10 +0000 Subject: [PATCH 04/11] Simplify layer artifact reuse --- lib/images/layer_artifact.go | 52 ++++++++++++++++++++++-------------- 1 file changed, 32 insertions(+), 20 deletions(-) diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index 83ef0a88c..6bce77699 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -110,21 +110,15 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { // and installed atomically; an interrupted build leaves only temp files that // the next attempt replaces. 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) + layerHex, err := layerDigestHex(desc) + if err != nil { + return nil, err } m.layerLifecycleMu.Lock() defer m.layerLifecycleMu.Unlock() - - if record, err := readLayerRecord(m.paths, layerHex); 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. + if record, found, err := m.existingLayerArtifact(layerHex, desc); found || err != nil { + return record, err } blobPath := m.paths.OCICacheBlob(layerHex) @@ -157,9 +151,9 @@ func (m *manager) materializeLayerArtifact(desc layerDescriptor) (*layerArtifact } func (m *manager) materializeLayerArtifactFromTree(desc layerDescriptor, tree layerTree) (*layerArtifact, error) { - layerHex := strings.TrimPrefix(desc.Digest, "sha256:") - if layerHex == "" || strings.Contains(layerHex, "/") || layerHex == "." || strings.Contains(layerHex, "..") { - return nil, fmt.Errorf("invalid layer digest: %s", desc.Digest) + layerHex, err := layerDigestHex(desc) + if err != nil { + return nil, err } if tree.path == "" || tree.stats == nil { return nil, fmt.Errorf("missing layer staging tree: %s", desc.Digest) @@ -170,16 +164,34 @@ func (m *manager) materializeLayerArtifactFromTree(desc layerDescriptor, tree la m.layerLifecycleMu.Lock() defer m.layerLifecycleMu.Unlock() - if record, err := readLayerRecord(m.paths, layerHex); 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 - } + if record, found, err := m.existingLayerArtifact(layerHex, desc); found || err != nil { + return record, err } return m.installLayerArtifact(desc, layerHex, tree.path, tree.stats) } +func layerDigestHex(desc layerDescriptor) (string, 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) + } + return layerHex, nil +} + +func (m *manager) existingLayerArtifact(layerHex string, desc layerDescriptor) (*layerArtifact, bool, error) { + record, err := readLayerRecord(m.paths, layerHex) + if err != nil { + return nil, true, err + } + if record == nil || !record.matches(desc) { + return nil, false, nil + } + if _, err := os.Stat(layerArtifactPath(m.paths, layerHex)); err == nil { + return record, true, nil + } + return nil, false, nil +} + func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir string, stats *unpackStats) (*layerArtifact, error) { record := &layerArtifact{ SchemaVersion: layerRecordSchemaVersion, From a82fc1c63873bd25526f1c5d30b28857d5e03f2f Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 18:59:55 +0000 Subject: [PATCH 05/11] Simplify image reference resolution --- lib/images/manager.go | 197 ++++++++++++++++++++++-------------------- 1 file changed, 105 insertions(+), 92 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index 9d67d4889..e9d631907 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -224,49 +224,9 @@ func (m *manager) CreateImage(ctx context.Context, req CreateImageRequest) (*Ima m.createMu.Lock() defer m.createMu.Unlock() - // Check if we already have this digest (deduplication) - if meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex()); err == nil { - // Don't cache failed builds - allow retry - if meta.Status == StatusFailed { - // Remove a failed tag before retrying. A tag for another digest is - // intentionally left untouched until the replacement is ready. - _ = deleteTagsForDigest(m.paths, ref.Repository(), ref.DigestHex()) - // Clean up the failed build directory so we can retry. Shared content - // is retained if another tag or pull still references it. - if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { - return nil, fmt.Errorf("remove failed image: %w", err) - } - // Fall through to re-queue the build - } else { - // Keep an existing ready tag visible until a replacement is ready. A - // pending tag is only created when it does not already point elsewhere. - if ref.Tag() != "" { - var tagErr error - if meta.Status == StatusReady { - tagErr = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } else { - tagErr = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) - } - if tagErr != nil { - return nil, fmt.Errorf("create image tag: %w", tagErr) - } - } - img := meta.toImage() - img.Name = ref.String() - if meta.Status == StatusReady { - return img, nil - } - if !m.inflightCredentialsMatch(ref.Digest(), req.Credentials) { - return nil, fmt.Errorf("%w: retry after the current pull completes", ErrCredentialConflict) - } - if meta.Status == StatusPending { - img.QueuePosition = m.queue.GetPosition(meta.Digest) - } - return img, nil - } + if img, found, err := m.reuseExistingImage(ref, req.Credentials); found || err != nil { + return img, err } - - // Don't have this digest yet, queue the build return m.createAndQueueImage(ref, req, platform) } @@ -330,6 +290,42 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest return m.createAndQueueImage(ref, CreateImageRequest{Name: imageRef}, hostPlatform()) } +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 + } + if meta.Status == StatusFailed { + _ = deleteTagsForDigest(m.paths, ref.Repository(), ref.DigestHex()) + if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), ref.DigestHex(), false); err != nil { + return nil, true, fmt.Errorf("remove failed image: %w", err) + } + return nil, false, nil + } + if ref.Tag() != "" { + if meta.Status == StatusReady { + err = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } else { + err = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + } + if err != nil { + return nil, true, fmt.Errorf("create image tag: %w", err) + } + } + img := meta.toImage() + img.Name = 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) inflightCredentialsMatch(digest string, credentials *authn.AuthConfig) bool { inflight := m.inflightPulls[digest] var existingFingerprint [32]byte @@ -853,20 +849,13 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, m.createMu.Lock() defer m.createMu.Unlock() - digestHex := sourceRef.DigestHex() - if !sourceRef.IsDigest() { - digestHex, err = resolveTag(m.paths, sourceRef.Repository(), sourceRef.Tag()) - if err != nil { - return nil, err - } - } - - meta, err := readMetadata(m.paths, sourceRef.Repository(), digestHex) + digestHex, err := m.resolveTagSource(sourceRef) if err != nil { return nil, err } - if meta.Status != StatusReady { - return nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) + meta, err := m.readyImageMetadata(sourceRef.Repository(), digestHex) + if err != nil { + return nil, err } previousDigest := "" @@ -894,6 +883,24 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, return img, nil } +func (m *manager) resolveTagSource(ref *NormalizedRef) (string, error) { + if ref.IsDigest() { + return ref.DigestHex(), nil + } + return resolveTag(m.paths, ref.Repository(), ref.Tag()) +} + +func (m *manager) readyImageMetadata(repository, digestHex string) (*imageMetadata, error) { + meta, err := readMetadata(m.paths, repository, digestHex) + if err != nil { + return nil, err + } + if meta.Status != StatusReady { + return nil, fmt.Errorf("%w: %s", ErrImageNotReady, meta.Status) + } + return meta, nil +} + func (m *manager) DeleteImage(ctx context.Context, name string) error { // Parse and normalize the reference ref, err := ParseNormalizedRef(name) @@ -1015,41 +1022,12 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { if err != nil { return fmt.Errorf("parse image name: %w", err) } - - // Wait for the image to appear in the store. In the build flow, the - // registry triggers ImportLocalImage asynchronously after a push, so the - // image may not exist when the build manager calls WaitForReady. - const maxWaitForExist = 30 * time.Second - const pollInterval = 100 * time.Millisecond - - var img *Image - deadline := time.Now().Add(maxWaitForExist) - for { - if !ref.IsDigest() { - img = m.findRequestedTagImage(ref) - } - if img == nil { - img, err = m.GetImage(ctx, name) - } - if img != nil { - break - } - if time.Now().After(deadline) { - return fmt.Errorf("get image: %w", err) - } - select { - case <-ctx.Done(): - return ctx.Err() - case <-time.After(pollInterval): - } + img, err := m.waitForImage(ctx, name, ref) + if err != nil { + return err } - - // Check if already in terminal state - switch img.Status { - case StatusReady: - return nil - case StatusFailed: - return conversionFailedErr(img.Error, nil) + if terminal, err := terminalImageError(img); terminal { + return err } digestHex := strings.TrimPrefix(img.Digest, "sha256:") @@ -1062,11 +1040,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 { - switch img.Status { - case StatusReady: - return nil - case StatusFailed: - return conversionFailedErr(img.Error, nil) + if terminal, terminalErr := terminalImageError(img); terminal { + return terminalErr } } @@ -1082,6 +1057,44 @@ func (m *manager) WaitForReady(ctx context.Context, name string) error { } } +func (m *manager) waitForImage(ctx context.Context, name string, ref *NormalizedRef) (*Image, error) { + const maxWaitForExist = 30 * time.Second + const pollInterval = 100 * time.Millisecond + var lastErr error + deadline := time.Now().Add(maxWaitForExist) + for { + var img *Image + if !ref.IsDigest() { + img = m.findRequestedTagImage(ref) + } + if img == nil { + img, lastErr = m.GetImage(ctx, name) + } + if img != nil { + return img, nil + } + if time.Now().After(deadline) { + return nil, fmt.Errorf("get image: %w", lastErr) + } + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(pollInterval): + } + } +} + +func 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) From 9928612bbc44217e0f7167e2b417d19487e6a5b1 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:00:28 +0000 Subject: [PATCH 06/11] Simplify image tag installation --- lib/images/manager.go | 21 +++++++++++++++------ 1 file changed, 15 insertions(+), 6 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index e9d631907..2f8faa99b 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -865,12 +865,8 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, return nil, tagErr } - if sourceRef.Repository() != targetRef.Repository() { - if err := promoteImageToContent(m.paths, sourceRef.Repository(), digestHex, meta, targetRef.Repository(), targetRef.Tag()); err != nil { - return nil, fmt.Errorf("create image alias: %w", err) - } - } else 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 { @@ -883,6 +879,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, targetRef.Repository(), 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 From 9a4cc427d1122f047e154aaa6219f5b8f425f63e Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:00:53 +0000 Subject: [PATCH 07/11] Simplify ready image finalization --- lib/images/manager.go | 23 ++++++++++++++--------- 1 file changed, 14 insertions(+), 9 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index 2f8faa99b..31b03daac 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -667,19 +667,24 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i } m.notifyReady(ref.DigestHex(), StatusReady, nil) - if ref.Tag() != "" { - currentDigest, tagErr := resolveTag(m.paths, ref.Repository(), ref.Tag()) - shouldPublish := tagErr == nil && (currentDigest == ref.DigestHex() || currentDigest == meta.PreviousTagDigest) - if shouldPublish { - if err := createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()); err != nil { - fmt.Fprintf(os.Stderr, "Warning: failed to create tag symlink: %v\n", err) - } - } - } + m.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) + if err != nil || (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" From a4135b2d0550aa1ef97859529e8eb99cda600034 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:30:44 +0000 Subject: [PATCH 08/11] Harden lifecycle failure handling --- lib/images/manager.go | 40 ++++++++++++++++++++++++++++++++++------ lib/images/oci_public.go | 5 ++++- lib/images/storage.go | 25 ++++++++++++++++++++++++- 3 files changed, 62 insertions(+), 8 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index 31b03daac..de845feb6 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -77,6 +77,7 @@ type manager struct { createMu sync.Mutex diskUsageMu sync.RWMutex layerLifecycleMu sync.Mutex + tagGenerations map[string]uint64 layerRefMu sync.Mutex inflightLayerRefs map[string]int diskUsageLoaded bool @@ -108,6 +109,7 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma inflightPulls: make(map[string]*inflightImagePull), borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, layerEvictionGrace: layerEvictionGracePeriod, + tagGenerations: make(map[string]uint64), inflightLayerRefs: make(map[string]int), readySubscribers: make(map[string][]chan StatusEvent), } @@ -326,6 +328,16 @@ func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthCo return img, true, nil } +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) inflightCredentialsMatch(digest string, credentials *authn.AuthConfig) bool { inflight := m.inflightPulls[digest] var existingFingerprint [32]byte @@ -415,6 +427,10 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, } } + tagGeneration := uint64(0) + if ref.Tag() != "" { + tagGeneration = m.nextTagGeneration(ref.Repository(), ref.Tag()) + } meta := &imageMetadata{ Name: ref.String(), Digest: ref.Digest(), @@ -426,6 +442,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, Tags: tags.Clone(req.Tags), RequestedTag: ref.Tag(), PreviousTagDigest: previousTagDigest, + TagGeneration: tagGeneration, CreatedAt: time.Now(), } @@ -625,11 +642,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 @@ -662,7 +674,15 @@ 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) } @@ -677,7 +697,8 @@ func (m *manager) publishReadyTag(ref *ResolvedRef, meta *imageMetadata) { return } currentDigest, err := resolveTag(m.paths, ref.Repository(), meta.RequestedTag) - if err != nil || (currentDigest != ref.DigestHex() && currentDigest != meta.PreviousTagDigest) { + 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 { @@ -741,6 +762,9 @@ func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err erro } 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 } @@ -853,6 +877,7 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, m.createMu.Lock() defer m.createMu.Unlock() + m.nextTagGeneration(targetRef.Repository(), targetRef.Tag()) digestHex, err := m.resolveTagSource(sourceRef) if err != nil { @@ -877,6 +902,8 @@ func (m *manager) TagImage(ctx context.Context, source, target string) (*Image, if err := removeDigestIfUnreferenced(m.paths, targetRef.Repository(), previousDigest, true); err != nil { return nil, fmt.Errorf("remove replaced image: %w", err) } + m.evictUnreferencedLayerArtifacts() + m.refreshDiskUsageTotals() } img := meta.toImage() @@ -946,6 +973,7 @@ func (m *manager) DeleteImage(ctx context.Context, name string) error { } tag := ref.Tag() + m.nextTagGeneration(repository, tag) // Resolve the tag to get the digest before deleting digestHex, err := resolveTag(m.paths, repository, tag) diff --git a/lib/images/oci_public.go b/lib/images/oci_public.go index 7d336745b..6b5500ca8 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 cleanupLayerTrees(result.LayerTrees) + } if err != nil { return fmt.Errorf("pull and unpack: %w", err) } diff --git a/lib/images/storage.go b/lib/images/storage.go index 7e060d9ed..cc2897963 100644 --- a/lib/images/storage.go +++ b/lib/images/storage.go @@ -233,8 +233,17 @@ func readMetadataAt(layout imageLayout) (*imageMetadata, error) { return &meta, nil } -func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata, target ...string) error { +func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, sourceMeta *imageMetadata, target ...string) (retErr error) { contentReady := false + var installedTarget *stagedTagSymlink + defer func() { + if retErr == nil || installedTarget == nil { + return + } + if restoreErr := restoreTagSymlink(installedTarget); restoreErr != nil { + retErr = errors.Join(retErr, fmt.Errorf("restore target tag: %w", restoreErr)) + } + }() if contentMeta, err := readContentMetadata(p, digestHex); err == nil { contentReady = contentMeta.Status == StatusReady } @@ -273,6 +282,7 @@ func promoteImageToContent(p *paths.Paths, sourceRepository, digestHex string, s _ = os.RemoveAll(staged.tempDir) return fmt.Errorf("install target tag: %w", err) } + installedTarget = &staged if err := removeStaleTagSymlink(p, &staged); err != nil { return fmt.Errorf("remove stale target tag: %w", err) } @@ -399,6 +409,19 @@ func stageTagSymlink(p *paths.Paths, repository, tag, digestHex string) (stagedT }, nil } +func restoreTagSymlink(staged *stagedTagSymlink) error { + if err := os.Remove(staged.linkPath); err != nil && !os.IsNotExist(err) { + return err + } + if !staged.previous.exists { + return nil + } + if err := os.MkdirAll(filepath.Dir(staged.linkPath), 0755); err != nil { + return err + } + return os.Symlink(staged.previous.target, staged.linkPath) +} + func readSymlinkState(path string) (symlinkState, error) { target, err := os.Readlink(path) if err != nil { From 16859e7953e129251f5ed3c863ff027fdd4c8efb Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:39:21 +0000 Subject: [PATCH 09/11] Fix lifecycle integrity and cleanup --- lib/images/layer_artifact.go | 99 ++++++++++++++++++++++++++++++++---- 1 file changed, 88 insertions(+), 11 deletions(-) diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index 6bce77699..5dc7ac0ce 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -63,7 +63,10 @@ type whiteoutRecord struct { } func (a *layerArtifact) matches(desc layerDescriptor) bool { - return a.Digest == desc.Digest && a.Format == layerArtifactFormat() + if a.Digest != desc.Digest || a.Format != layerArtifactFormat() { + return false + } + return desc.DiffID == "" || a.DiffID == desc.DiffID } const ( @@ -101,6 +104,9 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { if err := json.Unmarshal(data, &record); err != nil { return nil, fmt.Errorf("unmarshal layer record: %w", err) } + if err := record.validate(); err != nil { + return nil, fmt.Errorf("invalid layer record: %w", err) + } return &record, nil } @@ -109,6 +115,25 @@ func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { // 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) { + return fmt.Errorf("invalid digest or format") + } + if a.SizeBytes < 0 || a.UnpackedBytes < 0 || a.Entries < 0 { + return fmt.Errorf("invalid size or entry counts") + } + if a.Format == layerFormatExt4 && a.Options.Compression != "" { + return fmt.Errorf("ext4 artifact has compression options") + } + if a.Format == layerFormatErofs && a.Options.Compression != "lz4" { + return fmt.Errorf("erofs artifact has invalid compression options") + } + return nil +} + func (m *manager) materializeLayerArtifact(desc layerDescriptor) (*layerArtifact, error) { layerHex, err := layerDigestHex(desc) if err != nil { @@ -198,7 +223,7 @@ func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir Digest: desc.Digest, DiffID: desc.DiffID, Format: layerArtifactFormat(), - Options: layerArtifactOptions{Compression: "lz4"}, + Options: artifactOptions(layerArtifactFormat()), UnpackedBytes: stats.unpackedBytes, Entries: stats.entries, Whiteouts: stats.whiteouts, @@ -234,6 +259,13 @@ 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 @@ -586,7 +618,10 @@ func applyLayerTree(layerDir, targetDir string) error { if err != nil { return err } - targetParent := filepath.Join(targetDir, filepath.Dir(rel)) + targetParent, err := safeJoin(targetDir, filepath.Dir(rel)) + if err != nil { + return err + } if base == opaqueWhiteout { return clearDirContents(targetParent) } @@ -646,7 +681,7 @@ func clearDirContents(dir string) error { return err } if !info.IsDir() { - return fmt.Errorf("opaque whiteout target is not a directory: %s", dir) + return removePath(dir) } entries, err := os.ReadDir(dir) if err != nil { @@ -680,7 +715,7 @@ func copyEntryInto(src, dst string, hardlinks map[hardlinkIdentity]string) error case 0: return copyRegularEntry(src, dst, info, hardlinks) case fs.ModeDir: - return copyDirectoryEntry(dst, info) + return copyDirectoryEntry(src, dst, info) case fs.ModeSymlink: return copySymlinkEntry(src, dst) default: @@ -702,10 +737,10 @@ func copyRegularEntry(src, dst string, info os.FileInfo, hardlinks map[hardlinkI if err := copyFileContents(src, dst); err != nil { return err } - return copyEntryMetadata(dst, info) + return copyEntryMetadata(src, dst, info) } -func copyDirectoryEntry(dst string, info os.FileInfo) error { +func copyDirectoryEntry(src, dst string, info os.FileInfo) error { if existing, err := os.Lstat(dst); err == nil && !existing.IsDir() { if err := removePath(dst); err != nil { return err @@ -714,7 +749,7 @@ func copyDirectoryEntry(dst string, info os.FileInfo) error { if err := os.MkdirAll(dst, info.Mode().Perm()); err != nil { return err } - return copyEntryMetadata(dst, info) + return copyEntryMetadata(src, dst, info) } func copySymlinkEntry(src, dst string) error { @@ -746,7 +781,7 @@ func copySpecialEntry(src, dst string, info os.FileInfo) error { if err := unix.Mknod(dst, mode|uint32(info.Mode().Perm()), int(stat.Rdev)); err != nil { return err } - return copyEntryMetadata(dst, info) + return copyEntryMetadata(src, dst, info) } func specialFileMode(mode fs.FileMode) (uint32, error) { @@ -762,19 +797,61 @@ func specialFileMode(mode fs.FileMode) (uint32, error) { } } -func copyEntryMetadata(dst string, info os.FileInfo) error { +func copyEntryMetadata(src, dst string, info os.FileInfo) error { if stat, ok := info.Sys().(*syscall.Stat_t); ok { if err := os.Lchown(dst, int(stat.Uid), int(stat.Gid)); err != nil && !errors.Is(err, os.ErrPermission) && !errors.Is(err, unix.EPERM) { return err } } if info.Mode()&os.ModeSymlink == 0 { - if err := os.Chmod(dst, info.Mode().Perm()); err != nil { + mode := info.Mode().Perm() | info.Mode()&(os.ModeSetuid|os.ModeSetgid|os.ModeSticky) + if err := os.Chmod(dst, mode); err != nil { return err } if err := os.Chtimes(dst, info.ModTime(), info.ModTime()); err != nil { return err } + if err := copyXattrs(src, dst); err != nil { + return err + } + } + return nil +} + +func copyXattrs(src, dst string) error { + size, err := unix.Llistxattr(src, nil) + if err != nil { + if errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EPERM) { + return nil + } + return err + } + names := make([]byte, size) + if size > 0 { + n, err := unix.Llistxattr(src, names) + if err != nil { + return err + } + names = names[:n] + } + for _, name := range strings.Split(strings.TrimSuffix(string(names), "\\x00"), "\\x00") { + if name == "" { + continue + } + size, err := unix.Lgetxattr(src, name, nil) + if err != nil { + if errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EPERM) || errors.Is(err, unix.ENODATA) { + continue + } + return err + } + value := make([]byte, size) + if _, err := unix.Lgetxattr(src, name, value); err != nil { + return err + } + if err := unix.Lsetxattr(dst, name, value, 0); err != nil && !errors.Is(err, unix.ENOTSUP) && !errors.Is(err, unix.EPERM) { + return err + } } return nil } From 4f9353ed85940fc43349d10a4f5d2b0509ba6d2e Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:41:57 +0000 Subject: [PATCH 10/11] Require valid models for new images --- lib/images/manager.go | 13 +++++++------ lib/images/recovery_regression_test.go | 2 +- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/lib/images/manager.go b/lib/images/manager.go index de845feb6..9e23f385d 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -656,12 +656,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 diff --git a/lib/images/recovery_regression_test.go b/lib/images/recovery_regression_test.go index b979939a9..ad3398316 100644 --- a/lib/images/recovery_regression_test.go +++ b/lib/images/recovery_regression_test.go @@ -65,7 +65,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 { From 24cac62bac6ad60c3bf8d0c1a67c5655553f9025 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:46:27 +0000 Subject: [PATCH 11/11] Address Bugbot image lifecycle findings --- lib/images/layer_artifact.go | 2 +- lib/images/manager.go | 38 +++++++++++++++++++++++++++++++++++- 2 files changed, 38 insertions(+), 2 deletions(-) diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index 5dc7ac0ce..50c5903f1 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -197,7 +197,7 @@ func (m *manager) materializeLayerArtifactFromTree(desc layerDescriptor, tree la func layerDigestHex(desc layerDescriptor) (string, error) { layerHex := strings.TrimPrefix(desc.Digest, "sha256:") - if layerHex == "" || strings.Contains(layerHex, "/") || layerHex == "." || strings.Contains(layerHex, "..") { + if err := paths.ValidatePathComponent(layerHex); err != nil { return "", fmt.Errorf("invalid layer digest: %s", desc.Digest) } return layerHex, nil diff --git a/lib/images/manager.go b/lib/images/manager.go index 9e23f385d..025fd0617 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -271,6 +271,7 @@ func (m *manager) ImportLocalImage(ctx context.Context, repo, reference, digest if ref.Tag() != "" { var tagErr error if meta.Status == StatusReady { + m.nextTagGeneration(ref.Repository(), ref.Tag()) tagErr = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) } else { tagErr = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) @@ -306,6 +307,7 @@ func (m *manager) reuseExistingImage(ref *ResolvedRef, credentials *authn.AuthCo } if ref.Tag() != "" { if meta.Status == StatusReady { + m.nextTagGeneration(ref.Repository(), ref.Tag()) err = createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) } else { err = ensurePendingTag(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) @@ -333,11 +335,42 @@ func tagGenerationKey(repository, tag string) string { } 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] } +func (m *manager) restoreTagGenerations(metas []*imageMetadata) { + m.createMu.Lock() + defer m.createMu.Unlock() + if m.tagGenerations == nil { + m.tagGenerations = make(map[string]uint64) + } + for _, meta := range metas { + if meta.RequestedTag == "" { + 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 + } + } +} + +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 @@ -563,7 +596,9 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials if meta.Status == StatusReady { // Another build completed first; last-pull-wins repoints the tag. if ref.Tag() != "" { - createTagSymlink(m.paths, ref.Repository(), ref.Tag(), ref.DigestHex()) + if err := m.claimReadyTag(ref); err != nil { + slog.Warn("failed to claim ready image tag", "repository", ref.Repository(), "tag", ref.Tag(), "error", err) + } } buildStatus = "success" return @@ -785,6 +820,7 @@ 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 {