Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions cmd/api/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,8 @@ type ociCacheGCRunner interface {

// compositeOCICacheRoots fans the GC's extra-root query out to every source
// that tracks cache blobs outside index.json: the embedded registry's
// BuildKit cache tags and the push manager's in-flight push digests.
// BuildKit cache tags, the push manager's in-flight push digests, and image
// metadata (manifest + layer digests of every non-failed image).
type compositeOCICacheRoots []ocicachegc.RootsProvider

func (c compositeOCICacheRoots) LiveCacheManifestDigests() []string {
Expand Down Expand Up @@ -598,7 +599,7 @@ func run() error {

ociGC, err := configureOCICacheGC(
app.Config,
compositeOCICacheRoots{app.Registry, app.PushManager},
compositeOCICacheRoots{app.Registry, app.PushManager, images.NewOCICacheRoots(paths.New(app.Config.DataDir))},
logger,
otelProvider.MeterFor(loglib.SubsystemImages),
otelProvider.TracerFor(loglib.SubsystemImages),
Expand Down
115 changes: 115 additions & 0 deletions lib/images/cache_roots.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
package images

import (
"sort"

"github.com/kernel/hypeman/lib/paths"
)

// LiveOCICacheDigests returns the content digests that image metadata keeps
// alive in the shared OCI cache: the manifest digest and every recorded layer
// digest of each non-failed image. Failed images contribute nothing so their
// blobs become collectable once nothing else roots them.
//
// In-flight builds are protected because pending/pulling/converting metadata
// already carries the manifest digest as soon as the build is queued, and
// layer digests recorded at finalize keep blobs alive even after the image is
// no longer rooted in the OCI layout index.
func LiveOCICacheDigests(p *paths.Paths) []string {
metas, err := listAllMetadata(p)
if err != nil {
return nil
}

seen := make(map[string]struct{})
out := make([]string, 0)
add := func(digest string) {
if digest == "" {
return
}
if _, ok := seen[digest]; ok {
return
}
seen[digest] = struct{}{}
out = append(out, digest)
}
for _, meta := range metas {
if meta.Status == StatusFailed {
continue
}
add(meta.Digest)
for _, layer := range meta.Layers {
add(layer.Digest)
}
}
sort.Strings(out)
return out
}

// OCICacheRoots exposes image metadata as extra roots for the OCI cache
// garbage collector (lib/ocicachegc). It satisfies that package's
// RootsProvider interface structurally so the images package does not need
// to import it.
type OCICacheRoots struct {
paths *paths.Paths
}

func NewOCICacheRoots(p *paths.Paths) OCICacheRoots {
return OCICacheRoots{paths: p}
}

// LiveCacheManifestDigests returns the manifest and layer digests kept alive
// by image metadata. Layer digests are not manifests; the collector marks
// them live and treats their blobs as opaque leaves.
func (r OCICacheRoots) LiveCacheManifestDigests() []string {
return LiveOCICacheDigests(r.paths)
}

// layerAccounting summarises the bytes of OCI layers referenced by image
// metadata. Each layer digest is counted once regardless of how many images
// reference it, so shared layers are not double-counted.
type layerAccounting struct {
// uniqueBytes counts every referenced layer once.
uniqueBytes int64
// sharedBytes is the subset of uniqueBytes referenced by more than one
// image.
sharedBytes int64
}

func computeLayerAccounting(metas []*imageMetadata) layerAccounting {
type layerEntry struct {
size int64
refs int
}
layers := make(map[string]*layerEntry)
for _, meta := range metas {
if meta.Status == StatusFailed {
continue
}
seenInImage := make(map[string]struct{}, len(meta.Layers))
for _, layer := range meta.Layers {
if layer.Digest == "" {
continue
}
if _, dup := seenInImage[layer.Digest]; dup {
continue
}
seenInImage[layer.Digest] = struct{}{}
entry, ok := layers[layer.Digest]
if !ok {
entry = &layerEntry{size: layer.Size}
layers[layer.Digest] = entry
}
entry.refs++
}
}

var accounting layerAccounting
for _, entry := range layers {
accounting.uniqueBytes += entry.size
if entry.refs > 1 {
accounting.sharedBytes += entry.size
}
}
return accounting
}
187 changes: 187 additions & 0 deletions lib/images/cache_roots_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,187 @@
package images

import (
"context"
"fmt"
"os"
"path/filepath"
"testing"
"time"

"github.com/kernel/hypeman/lib/ocicachegc"
"github.com/kernel/hypeman/lib/paths"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func seedMetadataWithLayers(t *testing.T, p *paths.Paths, imageRef, status string, layers []LayerRef) *imageMetadata {
t.Helper()

ref, err := ParseNormalizedRef(imageRef)
require.NoError(t, err)

meta := &imageMetadata{
Name: imageRef,
Digest: ref.Digest(),
Status: status,
BuildID: "build-" + ref.DigestHex(),
Layers: layers,
CreatedAt: time.Now(),
}
require.NoError(t, writeMetadata(p, ref.Repository(), ref.DigestHex(), meta))
if status == StatusReady {
// Ready metadata is only readable when its disk file exists.
layout := resolveImageLayout(p, ref.Repository(), ref.DigestHex())
require.NoError(t, os.MkdirAll(filepath.Dir(layout.disk), 0o755))
require.NoError(t, os.WriteFile(layout.disk, []byte("rootfs"), 0o644))
}
return meta
}

func TestLiveOCICacheDigestsProtectsReferencedAndInflightArtifacts(t *testing.T) {
p := paths.New(t.TempDir())

sharedLayer := LayerRef{Digest: "sha256:" + sha256Hash([]byte("shared-layer")), Size: 10}
readyOnlyLayer := LayerRef{Digest: "sha256:" + sha256Hash([]byte("ready-layer")), Size: 20}
failedLayer := LayerRef{Digest: "sha256:" + sha256Hash([]byte("failed-layer")), Size: 30}

ready := seedMetadataWithLayers(t, p,
"docker.io/library/alpine@sha256:"+sha256Hash([]byte("ready")),
StatusReady, []LayerRef{sharedLayer, readyOnlyLayer})
inflight := seedMetadataWithLayers(t, p,
"docker.io/library/nginx@sha256:"+sha256Hash([]byte("inflight")),
StatusPulling, []LayerRef{sharedLayer})
failed := seedMetadataWithLayers(t, p,
"docker.io/library/busybox@sha256:"+sha256Hash([]byte("failed")),
StatusFailed, []LayerRef{failedLayer})

got := LiveOCICacheDigests(p)

assert.Contains(t, got, ready.Digest)
assert.Contains(t, got, sharedLayer.Digest)
assert.Contains(t, got, readyOnlyLayer.Digest)
// In-flight pulls carry only the manifest digest; their layers are not
// recorded until finalize, but the manifest root protects blobs written
// so far.
assert.Contains(t, got, inflight.Digest)
// Failed images contribute nothing so their blobs stay collectable.
assert.NotContains(t, got, failed.Digest)
assert.NotContains(t, got, failedLayer.Digest)

// Shared digests appear once even though two images reference them.
count := 0
for _, digest := range got {
if digest == sharedLayer.Digest {
count++
}
}
assert.Equal(t, 1, count)
}

// TestLiveOCICacheDigestsLegacyMetadataWithoutLayers proves that metadata
// written before layer tracking (no "layers" key) still loads and protects
// its manifest digest.
func TestLiveOCICacheDigestsLegacyMetadataWithoutLayers(t *testing.T) {
p := paths.New(t.TempDir())

ref, err := ParseNormalizedRef("docker.io/library/alpine@sha256:" + sha256Hash([]byte("legacy")))
require.NoError(t, err)

// Raw legacy shape: no layers key, matching pre-change metadata.json files.
legacyMeta := fmt.Sprintf(`{
"name": %q,
"digest": %q,
"status": %q,
"size_bytes": 6,
"created_at": %q
}`, ref.String(), ref.Digest(), StatusReady, time.Now().Format(time.RFC3339Nano))
layout := resolveImageLayout(p, ref.Repository(), ref.DigestHex())
require.NoError(t, os.MkdirAll(filepath.Dir(layout.metadata), 0o755))
require.NoError(t, os.WriteFile(layout.metadata, []byte(legacyMeta), 0o644))
require.NoError(t, os.WriteFile(layout.disk, []byte("rootfs"), 0o644))

got := LiveOCICacheDigests(p)
assert.Equal(t, []string{ref.Digest()}, got)
}

func TestComputeLayerAccountingCountsSharedLayersOnce(t *testing.T) {
sharedLayer := LayerRef{Digest: "sha256:" + sha256Hash([]byte("shared")), Size: 100}
uniqueLayerA := LayerRef{Digest: "sha256:" + sha256Hash([]byte("a")), Size: 50}
uniqueLayerB := LayerRef{Digest: "sha256:" + sha256Hash([]byte("b")), Size: 25}

metas := []*imageMetadata{
{Status: StatusReady, Layers: []LayerRef{sharedLayer, uniqueLayerA, sharedLayer}},
{Status: StatusReady, Layers: []LayerRef{sharedLayer, uniqueLayerB}},
{Status: StatusFailed, Layers: []LayerRef{{Digest: "sha256:" + sha256Hash([]byte("failed")), Size: 999}}},
}

accounting := computeLayerAccounting(metas)

// 100 (shared) + 50 (a) + 25 (b); the duplicate sharedLayer within the
// first image is counted once.
assert.Equal(t, int64(175), accounting.uniqueBytes)
assert.Equal(t, int64(100), accounting.sharedBytes)
}

// TestOCICacheGCProtectsImageReferencedBlobs runs a real GC sweep against a
// cache where one blob is referenced only by image metadata (not by
// index.json) and another is unreferenced. The metadata-referenced blob must
// survive; the orphan must be collected.
func TestOCICacheGCProtectsImageReferencedBlobs(t *testing.T) {
p := paths.New(t.TempDir())

blobDir := p.OCICacheBlobDir()
require.NoError(t, os.MkdirAll(blobDir, 0o755))

referencedBlob := sha256Hash([]byte("referenced-layer"))
orphanBlob := sha256Hash([]byte("orphan-layer"))
for _, blob := range []string{referencedBlob, orphanBlob} {
require.NoError(t, os.WriteFile(filepath.Join(blobDir, blob), []byte("blob-"+blob), 0o644))
}

seedMetadataWithLayers(t, p,
"docker.io/library/alpine@sha256:"+sha256Hash([]byte("ready")),
StatusReady, []LayerRef{{Digest: "sha256:" + referencedBlob, Size: 4}})

collector, err := ocicachegc.NewCollector(p, time.Hour, 0, NewOCICacheRoots(p), nil, nil, nil)
require.NoError(t, err)

stats, err := collector.Collect(context.Background())
require.NoError(t, err)

assert.Equal(t, 1, stats.DeletedBlobs)
assert.Equal(t, int64(len("blob-"+orphanBlob)), stats.DeletedBytes)
assert.FileExists(t, filepath.Join(blobDir, referencedBlob))
assert.NoFileExists(t, filepath.Join(blobDir, orphanBlob))
}

func TestFinalizeImageRecordsLayers(t *testing.T) {
p := paths.New(t.TempDir())
m := &manager{paths: p}

digestStr := "sha256:" + sha256Hash([]byte("finalize"))
imageRef := "docker.io/library/alpine@" + digestStr
ref, err := ParseNormalizedRef(imageRef)
require.NoError(t, err)
resolved := NewResolvedRef(ref, digestStr)

require.NoError(t, writeMetadata(p, ref.Repository(), ref.DigestHex(), &imageMetadata{
Name: imageRef,
Digest: digestStr,
Status: StatusPending,
BuildID: "b1",
CreatedAt: time.Now(),
}))

diskTempPath := filepath.Join(t.TempDir(), "rootfs.tmp")
require.NoError(t, os.WriteFile(diskTempPath, []byte("disk"), 0o644))

layers := []LayerRef{{Digest: "sha256:" + sha256Hash([]byte("layer-1")), Size: 5}}
err = m.finalizeImage(resolved, &pullResult{Metadata: &containerMetadata{}, Layers: layers}, int64(len("disk")), "b1", diskTempPath)
require.NoError(t, err)

meta, err := readMetadata(p, ref.Repository(), ref.DigestHex())
require.NoError(t, err)
require.Equal(t, StatusReady, meta.Status)
require.Equal(t, layers, meta.Layers)
}
5 changes: 3 additions & 2 deletions lib/images/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,7 @@ func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize i
meta.Env = result.Metadata.Env
meta.Labels = result.Metadata.Labels
meta.WorkingDir = result.Metadata.WorkingDir
meta.Layers = result.Layers

if err := writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta); err != nil {
return fmt.Errorf("write final metadata: %w", err)
Expand Down Expand Up @@ -578,11 +579,11 @@ func (m *manager) recordPullResultMetrics(ctx context.Context, digest string, re
m.recordImageBuildPhase(ctx, digest, phase.Phase, phase.Duration, phase.Status, cacheStatus)
}
if result.Metadata != nil {
m.recordOCIImageMetrics(ctx, result.LayerCount, result.CompressedBytes, cacheStatus)
m.recordOCIImageMetrics(ctx, len(result.Layers), result.CompressedBytes, cacheStatus)
slog.InfoContext(ctx, "OCI image inspected",
"digest", digest,
"cache_status", cacheStatus,
"layer_count", result.LayerCount,
"layer_count", len(result.Layers),
"compressed_bytes", result.CompressedBytes,
)
}
Expand Down
16 changes: 16 additions & 0 deletions lib/images/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,15 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) {
return nil, err
}

referencedLayerBytes, err := meter.Int64ObservableGauge(
"hypeman_images_referenced_layer_bytes",
metric.WithDescription("Bytes of OCI layers referenced by image metadata; shared layers counted once"),
metric.WithUnit("By"),
)
if err != nil {
return nil, err
}

_, err = meter.RegisterCallback(
func(ctx context.Context, o metric.Observer) error {
// Report queue length
Expand All @@ -113,10 +122,17 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) {
o.ObserveInt64(imagesTotal, count,
metric.WithAttributes(attribute.String("status", status)))
}

accounting := computeLayerAccounting(metas)
o.ObserveInt64(referencedLayerBytes, accounting.uniqueBytes,
metric.WithAttributes(attribute.String("scope", "unique")))
o.ObserveInt64(referencedLayerBytes, accounting.sharedBytes,
metric.WithAttributes(attribute.String("scope", "shared")))
return nil
},
buildQueueLength,
imagesTotal,
referencedLayerBytes,
)
if err != nil {
return nil, err
Expand Down
Loading
Loading