diff --git a/lib/images/compose.go b/lib/images/compose.go index 9df8c6175..a799d240c 100644 --- a/lib/images/compose.go +++ b/lib/images/compose.go @@ -15,8 +15,9 @@ import ( // concurrently, and a failure between the remove and the rename leaves dest // absent. The export root is always 0755 regardless of the last layer's tar // root entry, matching the mode the previous unpack path created. A crash -// can also strand .compose-* staging directories in dest's parent, the same -// way .unpack-* directories can strand under layer builds. +// can also strand .compose-* staging directories in dest's parent build +// directory; the next compose attempt for the same digest removes stale +// ones before creating its own. func (c *ociClient) composeRootfs(ctx context.Context, dest, layoutTag string, model *imageManifestModel) error { if err := validateManifestModel(layoutTag, model); err != nil { return fmt.Errorf("validate manifest model: %w", err) @@ -25,6 +26,14 @@ func (c *ociClient) composeRootfs(ctx context.Context, dest, layoutTag string, m if err := os.MkdirAll(parent, 0755); err != nil { return fmt.Errorf("create compose parent: %w", err) } + // The build directory is digest-keyed, so any leftover .compose-* sibling + // is garbage from a crashed build of the same digest. + leftovers, _ := filepath.Glob(filepath.Join(parent, ".compose-*")) + for _, leftover := range leftovers { + if err := removePath(leftover); err != nil { + slog.Warn("failed to remove stale compose staging directory", "dir", leftover, "error", err) + } + } staging, err := os.MkdirTemp(parent, ".compose-*") if err != nil { return fmt.Errorf("create compose directory: %w", err) diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go index e2faf8093..11e401789 100644 --- a/lib/images/layer_artifact.go +++ b/lib/images/layer_artifact.go @@ -154,9 +154,8 @@ func discardLayerCache(p *paths.Paths, layerHex string) error { // The layer is unpacked into an isolated temp directory, converted to the // default image format, and installed atomically. Normal failures remove the // temp directory; a crash mid-build can leave a stale .unpack-* directory -// behind, which reconciliation landing with the pull integration is expected -// to sweep. No production caller yet: pull integration and -// composition land in later changes. +// behind, which the startup sweep in layer_gc.go removes once it ages past +// the eviction grace period. // // Concurrent callers share one build. The build itself is detached from the // initiating caller's cancellation so one cancelled pull cannot fail every diff --git a/lib/images/layer_artifact_test.go b/lib/images/layer_artifact_test.go index 28e3cca7e..f0d7f6f0c 100644 --- a/lib/images/layer_artifact_test.go +++ b/lib/images/layer_artifact_test.go @@ -32,17 +32,19 @@ const whiteoutPrefix = ".wh." const testTarGzMediaType = "application/vnd.oci.image.layer.v1.tar+gzip" -// writeLayerTestLayout writes img into the shared OCI cache of p tagged with -// the image's digest, mirroring pullToOCILayout. -func writeLayerTestLayout(t *testing.T, p *paths.Paths, img gcr.Image) { +// writeLayerTestLayout writes images into the shared OCI cache of p tagged +// with each image's digest, mirroring pullToOCILayout. +func writeLayerTestLayout(t *testing.T, p *paths.Paths, imgs ...gcr.Image) { t.Helper() - digest, err := img.Digest() - require.NoError(t, err) layoutPath, err := layout.Write(p.SystemOCICache(), empty.Index) require.NoError(t, err) - require.NoError(t, layoutPath.AppendImage(img, layout.WithAnnotations(map[string]string{ - "org.opencontainers.image.ref.name": digestToLayoutTag(digest.String()), - }))) + 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()), + }))) + } } func layerDescFromImage(t *testing.T, img gcr.Image, index int) layerDescriptor { diff --git a/lib/images/layer_gc.go b/lib/images/layer_gc.go new file mode 100644 index 000000000..9f7a071a9 --- /dev/null +++ b/lib/images/layer_gc.go @@ -0,0 +1,232 @@ +package images + +import ( + "context" + "fmt" + "io/fs" + "log/slog" + "os" + "path/filepath" + "strings" + "sync" + "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 images tree — both content and +// legacy layouts write their model as a manifest.json with the digest as its +// parent directory — plus the digests currently referenced by in-flight +// builds. Layer artifacts in this set are protected from eviction. Unreadable +// manifest models are skipped with a warning so one corrupt record cannot +// disable eviction entirely. +// +// Callers must hold createMu so the in-flight map read is ordered with model +// writes during finalization. +func (m *manager) referencedLayerDigests() (map[string]struct{}, error) { + refs := make(map[string]struct{}, len(m.inflightLayerRefs)) + for digestHex := range m.inflightLayerRefs { + refs[digestHex] = struct{}{} + } + err := filepath.WalkDir(m.paths.ImagesDir(), func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + // An incomplete walk means an incomplete reference set; the caller + // must not evict against it. + return err + } + if entry.IsDir() || entry.Name() != "manifest.json" { + return nil + } + digestHex := filepath.Base(filepath.Dir(path)) + model, readErr := readManifestModelAt(path, digestHex) + if readErr != nil { + slog.Warn("skipping unreadable manifest model for layer eviction", "digest", digestHex, "error", readErr) + return nil + } + if model == nil { + return nil + } + for _, layer := range model.Layers { + refs[strings.TrimPrefix(layer.Digest, "sha256:")] = struct{}{} + } + return nil + }) + if err != nil { + return nil, fmt.Errorf("walk image manifests: %w", err) + } + return refs, nil +} + +// inflightLayerRef is the handle returned by retainInflightLayers. Its +// release is idempotent: finalization releases the refs as soon as the +// manifest model is durable, and the build's deferred release becomes a +// no-op afterwards. +type inflightLayerRef struct { + once sync.Once + digestHexes []string +} + +func (r *inflightLayerRef) release(m *manager) { + r.once.Do(func() { + m.createMu.Lock() + defer m.createMu.Unlock() + m.releaseInflightLayerRefsLocked(r.digestHexes) + }) +} + +// releaseLocked is release for callers already holding createMu. +func (r *inflightLayerRef) releaseLocked(m *manager) { + r.once.Do(func() { m.releaseInflightLayerRefsLocked(r.digestHexes) }) +} + +// retainInflightLayers registers one in-flight reference per digest so +// reconciliation cannot evict layers a build is materializing. +func (m *manager) retainInflightLayers(digestHexes []string) *inflightLayerRef { + m.createMu.Lock() + for _, digestHex := range digestHexes { + m.inflightLayerRefs[digestHex]++ + } + m.createMu.Unlock() + return &inflightLayerRef{digestHexes: digestHexes} +} + +func (m *manager) releaseInflightLayerRefsLocked(digestHexes []string) { + for _, digestHex := range digestHexes { + if m.inflightLayerRefs[digestHex] <= 1 { + delete(m.inflightLayerRefs, digestHex) + } else { + m.inflightLayerRefs[digestHex]-- + } + } +} + +// reconcileLayerStore evicts unreferenced layer artifacts and refreshes the +// cached disk usage totals so accounting reflects the removals. +func (m *manager) reconcileLayerStore() { + m.createMu.Lock() + defer m.createMu.Unlock() + m.reconcileLayerStoreLocked() +} + +// reconcileLayerStoreLocked is used by lifecycle operations that already hold +// createMu. Serializing reconciliation with manifest finalization prevents an +// eviction scan from racing a newly committed layer reference. +func (m *manager) reconcileLayerStoreLocked() { + m.evictUnreferencedLayerArtifacts() + m.refreshDiskUsageTotals() +} + +// evictUnreferencedLayerArtifacts removes layer artifacts that no image +// manifest model references, deleting the digest directory entirely. Artifacts +// newer than the grace period are kept so in-flight builds never lose work. +func (m *manager) evictUnreferencedLayerArtifacts() { + refs, err := m.referencedLayerDigests() + if err != nil { + // Evicting against a truncated reference set would delete artifacts + // belonging to images the walk never reached. + slog.Warn("skipping layer eviction: incomplete reference scan", "error", err) + return + } + + layersDir := m.paths.ImageLayersDir() + entries, err := os.ReadDir(layersDir) + if err != nil { + if !os.IsNotExist(err) { + slog.Warn("layer eviction failed to list layer store", "error", err) + } + return + } + + cutoff := time.Now().Add(-m.layerEvictionGrace) + evicted := 0 + var evictedBytes int64 + for _, entry := range entries { + if !entry.IsDir() { + continue + } + digestHex := entry.Name() + if _, referenced := refs[digestHex]; referenced { + continue + } + dirPath := filepath.Join(layersDir, digestHex) + info, statErr := os.Stat(dirPath) + if statErr != nil || info.ModTime().After(cutoff) { + continue + } + size, err := dirSize(dirPath) + if err != nil { + slog.Warn("failed to measure layer artifact size", "digest", digestHex, "error", err) + } + // removePath clears read-only directories restored from layer + // metadata, which os.RemoveAll cannot unlink through. + if err := removePath(dirPath); err != nil { + slog.Warn("failed to evict unreferenced layer artifact", "digest", digestHex, "error", err) + continue + } + evicted++ + evictedBytes += size + } + if evicted > 0 { + slog.Info("evicted unreferenced layer artifacts", "count", evicted, "bytes", evictedBytes) + m.recordLayerArtifactsEvicted(context.Background(), int64(evicted)) + } +} + +// isStaleTempDirName reports whether a directory name matches the temp +// prefixes builds use for staging, installs, and tag promotion. +func isStaleTempDirName(name string) bool { + for _, prefix := range []string{".unpack-", ".install-", ".tag-stage-"} { + if strings.HasPrefix(name, prefix) { + return true + } + } + return false +} + +// cleanStaleImageTempDirs removes temp directories left behind by builds that +// were interrupted mid-install, mid-materialization, or mid-tag promotion. +// The walk covers the whole images tree: layer and content staging dirs plus +// .tag-stage-* dirs, which are created under images//. Only +// directories older than the grace period are removed so live builds are +// never disturbed. +// +// This must stay a startup-only sweep: a staging dir's own mtime only moves +// when its direct children change, so a deep extraction running longer than +// the grace period can look stale while actively writing. A periodic sweep +// would need a heartbeat or a live-build registry first. +func (m *manager) cleanStaleImageTempDirs() { + cutoff := time.Now().Add(-m.layerEvictionGrace) + err := filepath.WalkDir(m.paths.ImagesDir(), func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if !entry.IsDir() { + return nil + } + name := entry.Name() + if !isStaleTempDirName(name) { + return nil + } + if info, err := entry.Info(); err == nil && info.ModTime().Before(cutoff) { + // removePath clears the read-only directories umoci restores from + // layer metadata, which os.RemoveAll cannot unlink through. + if err := removePath(path); err != nil { + slog.Warn("failed to remove stale image temp dir", "dir", path, "error", err) + } + } + return fs.SkipDir + }) + if err != nil { + slog.Warn("failed to clean stale image temp dirs", "root", m.paths.ImagesDir(), "error", err) + } +} diff --git a/lib/images/lifecycle_test.go b/lib/images/lifecycle_test.go new file mode 100644 index 000000000..6ac8a06dc --- /dev/null +++ b/lib/images/lifecycle_test.go @@ -0,0 +1,205 @@ +package images + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + gcr "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/empty" + "github.com/google/go-containerregistry/pkg/v1/mutate" + "github.com/kernel/hypeman/lib/paths" + "github.com/stretchr/testify/require" +) + +func layerStoreHexes(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 +} + +func imageDigest(t *testing.T, img gcr.Image) string { + t.Helper() + digest, err := img.Digest() + require.NoError(t, err) + return digest.String() +} + +// importAndWait imports an image and blocks until its build reaches ready. +func importAndWait(t *testing.T, m *manager, ctx context.Context, repo, tag, digest string) { + t.Helper() + events := make(chan StatusEvent, 2) + layoutTag := digestToLayoutTag(digest) + m.subscribeToReady(layoutTag, events) + defer m.unsubscribeFromReady(layoutTag, events) + _, err := m.ImportLocalImage(ctx, repo, tag, digest) + require.NoError(t, err) + select { + case event := <-events: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatalf("image %s did not become ready", repo) + } +} + +// TestSharedLayersMaterializeOnceAndEvictWithReferences is the end-to-end +// lifecycle: two images share a base layer, the shared artifact is created +// once, survives the deletion of one image, and is evicted only when its last +// reference is gone. +func TestSharedLayersMaterializeOnceAndEvictWithReferences(t *testing.T) { + if _, err := exec.LookPath("mkfs.erofs"); err != nil { + t.Skip("mkfs.erofs not available") + } + dataDir := t.TempDir() + p := paths.New(dataDir) + mgr, err := NewManager(p, 1, nil) + require.NoError(t, err) + m := mgr.(*manager) + m.layerEvictionGrace = 0 + + base := syntheticLayer(t, "base.txt", "shared base content") + topA := syntheticLayer(t, "a.txt", "app A payload") + topB := syntheticLayer(t, "b.txt", "app B payload") + + imgA, err := mutate.AppendLayers(empty.Image, base, topA) + require.NoError(t, err) + imgB, err := mutate.AppendLayers(empty.Image, base, topB) + require.NoError(t, err) + + writeLayerTestLayout(t, p, imgA, imgB) + digestA, digestB := imageDigest(t, imgA), imageDigest(t, imgB) + + baseManifest, err := imgA.Manifest() + require.NoError(t, err) + baseHex, topAHex := baseManifest.Layers[0].Digest.Hex, 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" + + importAndWait(t, m, ctx, repoA, "v1", digestA) + importAndWait(t, m, ctx, repoB, "v1", digestB) + + // The shared base layer materialized exactly once, alongside the two tops. + hexes := layerStoreHexes(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 = layerStoreHexes(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 = layerStoreHexes(t, p) + require.Empty(t, hexes, "unreferenced layer artifacts must be evicted") +} + +func newLifecycleTestManager(p *paths.Paths) *manager { + return &manager{ + paths: p, + inflightPulls: make(map[string]*inflightImagePull), + inflightLayerRefs: make(map[string]int), + } +} + +// TestLayerArtifactsCountedInOneBucket verifies the accounting contract: ready +// image bytes and the OCI+layer cache total stay disjoint, so consumers summing +// TotalImageBytes and TotalOCICacheBytes count layer bytes exactly once. +func TestLayerArtifactsCountedInOneBucket(t *testing.T) { + p := paths.New(t.TempDir()) + m := newLifecycleTestManager(p) + + digestHex := "cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01cd01" + require.NoError(t, os.MkdirAll(p.ImageLayerDir(digestHex), 0o755)) + payload := make([]byte, 4096) + require.NoError(t, os.WriteFile(p.ImageLayerArtifactForFormat(digestHex, string(DefaultImageFormat)), payload, 0o644)) + + readyBytes, cacheBytes, err := m.getDiskUsageTotals() + require.NoError(t, err) + require.GreaterOrEqual(t, cacheBytes, int64(len(payload))) + + totalBytes, err := m.TotalImageBytes(context.Background()) + require.NoError(t, err) + require.Equal(t, readyBytes, totalBytes) + + cacheTotal, err := m.TotalOCICacheBytes(context.Background()) + require.NoError(t, err) + require.Equal(t, cacheBytes, cacheTotal) +} + +func TestCleanStaleImageTempDirsRemovesOnlyOldDirectories(t *testing.T) { + p := paths.New(t.TempDir()) + m := newLifecycleTestManager(p) + m.layerEvictionGrace = time.Hour + + layersDir := p.ImageLayersDir() + staleDir := filepath.Join(layersDir, "ab12", ".unpack-stale") + freshDir := filepath.Join(layersDir, "cd34", ".unpack-fresh") + require.NoError(t, os.MkdirAll(staleDir, 0o755)) + require.NoError(t, os.MkdirAll(freshDir, 0o755)) + old := time.Now().Add(-2 * time.Hour) + require.NoError(t, os.Chtimes(staleDir, old, old)) + + m.cleanStaleImageTempDirs() + + _, err := os.Stat(staleDir) + require.True(t, os.IsNotExist(err), "stale temp dir must be removed") + _, err = os.Stat(freshDir) + require.NoError(t, err, "fresh temp dir must survive cleanup") +} + +func TestEvictionKeepsReferencedAndFreshArtifacts(t *testing.T) { + p := paths.New(t.TempDir()) + m := newLifecycleTestManager(p) + m.layerEvictionGrace = time.Hour + + referencedHex := "ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01ef01" + orphanFreshHex := "ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23ab23" + + // A manifest model referencing one layer protects it regardless of age. + model := &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: "sha256:" + referencedHex, + RootFSType: "layers", + Config: manifestConfigRef{ + Digest: "sha256:" + strings.Repeat("c", 64), + MediaType: "application/vnd.oci.image.config.v1+json", + DiffIDs: []string{"sha256:" + referencedHex}, + }, + Layers: []layerDescriptor{{Digest: "sha256:" + referencedHex, DiffID: "sha256:" + referencedHex}}, + } + require.NoError(t, writeManifestModel(p, referencedHex, model)) + require.NoError(t, os.MkdirAll(p.ImageLayerDir(referencedHex), 0o755)) + require.NoError(t, os.WriteFile(p.ImageLayerArtifactForFormat(referencedHex, string(DefaultImageFormat)), []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.ImageLayerArtifactForFormat(orphanFreshHex, string(DefaultImageFormat)), []byte("fresh"), 0o644)) + + m.reconcileLayerStore() + + hexes := layerStoreHexes(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 428e143c5..66bb149bd 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -77,6 +77,8 @@ type manager struct { queue *queue.Queue createMu sync.Mutex layerFlights singleflight.Group + inflightLayerRefs map[string]int + layerEvictionGrace time.Duration diskUsageMu sync.RWMutex tagGenerations map[string]uint64 requestedTags map[string]string // newest pull's digest per requested tag @@ -105,6 +107,8 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma ociClient: ociClient, queue: queue.New(maxConcurrentBuilds), inflightPulls: make(map[string]*inflightImagePull), + inflightLayerRefs: make(map[string]int), + layerEvictionGrace: layerEvictionGracePeriod, borrowedCredentialsTimeout: DefaultBorrowedCredentialsTimeout, readySubscribers: make(map[string][]chan StatusEvent), tagGenerations: make(map[string]uint64), @@ -121,6 +125,12 @@ func NewManager(p *paths.Paths, maxConcurrentBuilds int, meter metric.Meter) (Ma } m.RecoverInterruptedBuilds() + // Sweep temp dirs and evict layer orphans an unclean shutdown may have + // left behind. Recovered builds re-enqueued above run concurrently; they + // re-materialize from the blob cache if this pass evicts their previous + // attempt's unreferenced artifacts. + m.cleanStaleImageTempDirs() + m.reconcileLayerStore() // Keep legacy images readable in their existing layout and promote them only // when an operation needs shared content, such as a cross-repository tag. // Avoiding a startup-wide migration keeps startup bounded and independent of @@ -494,6 +504,22 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } m.recordPullMetrics(ctx, "success") + materialized, err := m.materializeLayerArtifacts(ctx, result) + if materialized != nil { + // Hold the in-flight references until the manifest model protecting + // the layers is durable: on the cache-hit path the artifacts' mtimes + // are too old for the eviction grace period to cover the gap. + // finalizeImage releases them under createMu once the model is + // written; this defer covers every path that returns before that. + defer materialized.release(m) + } + if err != nil { + // The rootfs is already composed from blobs, so materialization is + // best effort: log and continue with degraded sharing. + slog.Warn("layer materialization failed; continuing without shared artifacts", + "digest", ref.DigestHex(), "error", err) + } + // 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 { @@ -528,7 +554,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } finalizeStart := time.Now() - err = m.finalizeImage(ref, result, diskSize, buildID, diskTempPath) + err = m.finalizeImage(ref, result, diskSize, buildID, diskTempPath, materialized) m.recordImageBuildPhase(ctx, ref.Digest(), "finalize", time.Since(finalizeStart), phaseStatus(err), "not_applicable") if err != nil { if errors.Is(err, errStaleBuild) { @@ -541,7 +567,26 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials buildStatus = "success" } -func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string) error { +func (m *manager) materializeLayerArtifacts(ctx context.Context, result *pullResult) (*inflightLayerRef, error) { + if result.Manifest == nil || layerArtifactFormat() == "" { + return nil, nil + } + digestHexes := make([]string, 0, len(result.Manifest.Layers)) + for _, desc := range result.Manifest.Layers { + digestHexes = append(digestHexes, strings.TrimPrefix(desc.Digest, "sha256:")) + } + handle := m.retainInflightLayers(digestHexes) + for _, desc := range result.Manifest.Layers { + if _, err := m.materializeLayerArtifact(ctx, desc); err != nil { + // The handle's release is idempotent: the caller keeps the partial + // set protected until its path finishes with it. + return handle, err + } + } + return handle, nil +} + +func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string, materialized *inflightLayerRef) error { m.createMu.Lock() defer m.createMu.Unlock() @@ -584,6 +629,13 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i return rollbackFinalization(layout, modelPath, diskInstalled, modelWritten, fmt.Errorf("write manifest model: %w", err)) } modelWritten = true + // The model now protects these layers on disk, so the in-flight refs + // can go while still holding createMu: a delete racing the ready + // notification then reconciles against durable state. The build's + // deferred release is a no-op after this. + if materialized != nil { + materialized.releaseLocked(m) + } } meta.Status = StatusReady @@ -817,7 +869,7 @@ func (m *manager) deleteDigestImage(repository, digestHex string) error { return err } m.clearRequestedDigest(digestHex) - m.refreshDiskUsageTotals() + m.reconcileLayerStoreLocked() return nil } @@ -851,11 +903,13 @@ func (m *manager) deleteTaggedImage(repository, tag string) error { if err := removeDigestIfUnreferenced(m.paths, repository, digestHex, true); err != nil { return fmt.Errorf("delete orphaned digest %s: %w", digestHex, err) } - m.refreshDiskUsageTotals() + m.reconcileLayerStoreLocked() return nil } -// TotalImageBytes returns the total size of all ready images on disk. +// TotalImageBytes returns the total size of all ready images on disk. Shared +// layer artifacts are accounted separately via TotalOCICacheBytes so the two +// totals can be summed without double-counting. func (m *manager) TotalImageBytes(ctx context.Context) (int64, error) { readyImageBytes, _, err := m.getDiskUsageTotals() if err != nil { diff --git a/lib/images/manager_test.go b/lib/images/manager_test.go index 13a6835e2..6f24e6390 100644 --- a/lib/images/manager_test.go +++ b/lib/images/manager_test.go @@ -736,7 +736,7 @@ func TestDeleteAndRecreateDuringBuildTail(t *testing.T) { m.updateStatusByDigest(staleRef, StatusFailed, errors.New("stale build"), firstMeta.BuildID) staleBundle, err := m.ociClient.extractOCIImageBundle(digestHex) require.NoError(t, err) - require.ErrorIs(t, m.finalizeImage(staleRef, &pullResult{Metadata: staleBundle.Meta}, 1, firstMeta.BuildID, ""), errStaleBuild) + require.ErrorIs(t, m.finalizeImage(staleRef, &pullResult{Metadata: staleBundle.Meta}, 1, firstMeta.BuildID, "", nil), errStaleBuild) currentMeta, err = readMetadata(p, repo, digestHex) require.NoError(t, err) require.Equal(t, StatusPending, currentMeta.Status) diff --git a/lib/images/manifest_model.go b/lib/images/manifest_model.go index 16a913ce4..5cdfd8f30 100644 --- a/lib/images/manifest_model.go +++ b/lib/images/manifest_model.go @@ -149,7 +149,13 @@ func writeManifestModelAt(path, digestHex string, model *imageManifestModel) err // Missing models return (nil, nil): images converted before the manifest model // existed only have a flattened rootfs. func readManifestModel(p *paths.Paths, digestHex string) (*imageManifestModel, error) { - data, err := os.ReadFile(p.ImageContentManifestModel(digestHex)) + return readManifestModelAt(p.ImageContentManifestModel(digestHex), digestHex) +} + +// readManifestModelAt loads the manifest model at path, if present. Missing +// files return (nil, nil). +func readManifestModelAt(path, digestHex string) (*imageManifestModel, error) { + data, err := os.ReadFile(path) if err != nil { if os.IsNotExist(err) { return nil, nil diff --git a/lib/images/metrics.go b/lib/images/metrics.go index d860885b2..e58075905 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 } @@ -169,3 +179,11 @@ func (m *manager) recordOCIImageMetrics(ctx context.Context, layerCount int, com m.metrics.ociLayerCount.Record(ctx, int64(layerCount), attrs) m.metrics.ociCompressedBytes.Record(ctx, compressedBytes, attrs) } + +// recordLayerArtifactsEvicted counts layer artifacts removed by eviction. +func (m *manager) recordLayerArtifactsEvicted(ctx context.Context, evicted int64) { + if m.metrics == nil { + return + } + m.metrics.layerArtifactsEvicted.Add(ctx, evicted) +} diff --git a/lib/images/tag.go b/lib/images/tag.go index 302263d7e..c64415b43 100644 --- a/lib/images/tag.go +++ b/lib/images/tag.go @@ -92,7 +92,7 @@ func (m *manager) cleanupReplacedTag(ref *NormalizedRef, previousDigest, digestH if err := removeDigestIfUnreferenced(m.paths, ref.Repository(), previousDigest, true); err != nil { slog.Warn("failed to collect replaced image content", "repository", ref.Repository(), "digest", previousDigest, "error", err) } - m.refreshDiskUsageTotals() + m.reconcileLayerStoreLocked() } func parseTagReferences(source, target string) (*NormalizedRef, *NormalizedRef, error) { diff --git a/lib/resources/resource.go b/lib/resources/resource.go index 5943cb0bc..8af7e32d7 100644 --- a/lib/resources/resource.go +++ b/lib/resources/resource.go @@ -889,7 +889,7 @@ func (m *Manager) MaxImageStorageBytes() int64 { return int64(float64(capacity) * fraction) } -// CurrentImageStorageBytes returns the current image storage usage (OCI cache + rootfs). +// CurrentImageStorageBytes returns the current image storage usage (ready rootfs + OCI cache and shared layers). func (m *Manager) CurrentImageStorageBytes(ctx context.Context) (int64, error) { if m.imageLister == nil { return 0, nil