diff --git a/lib/images/compose.go b/lib/images/compose.go new file mode 100644 index 000000000..304732ca0 --- /dev/null +++ b/lib/images/compose.go @@ -0,0 +1,43 @@ +package images + +import ( + "context" + "fmt" + "os" + "path/filepath" +) + +// composeRootfs validates the persisted model and merges its 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. +func (c *ociClient) composeRootfs(dest, layoutTag string, model *imageManifestModel) error { + return c.composeRootfsContext(context.Background(), dest, layoutTag, model) +} + +func (c *ociClient) composeRootfsContext(ctx context.Context, dest, layoutTag string, model *imageManifestModel) error { + if err := validateManifestModel(layoutTag, model); err != nil { + return fmt.Errorf("validate manifest model: %w", err) + } + if err := os.MkdirAll(filepath.Dir(dest), 0755); err != nil { + return fmt.Errorf("create compose parent: %w", err) + } + staging, err := os.MkdirTemp(filepath.Dir(dest), ".compose-*") + if err != nil { + return fmt.Errorf("create compose directory: %w", err) + } + defer os.RemoveAll(staging) + + for i, desc := range model.Layers { + if _, err := unpackCachedLayer(ctx, c.cacheDir, desc, staging, composeOnDiskFormat()); err != nil { + return fmt.Errorf("apply layer %d (%s): %w", i, desc.Digest, err) + } + } + if err := os.RemoveAll(dest); err != nil { + return fmt.Errorf("replace compose directory: %w", err) + } + if err := os.Rename(staging, dest); err != nil { + return fmt.Errorf("install compose directory: %w", err) + } + return nil +} diff --git a/lib/images/compose_test.go b/lib/images/compose_test.go new file mode 100644 index 000000000..3e0dea77d --- /dev/null +++ b/lib/images/compose_test.go @@ -0,0 +1,207 @@ +package images + +import ( + "io" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + + 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" +) + +// composeTestImage builds the standard two-layer fixture: a base layer with +// content the top layer deletes, masks, replaces, and extends. +func composeTestImage(t *testing.T) gcr.Image { + t.Helper() + + base := specLayer(t, []tarEntrySpec{ + {name: "etc/", isDir: true, mode: 0755}, + {name: "etc/config.txt", content: "original", mode: 0644}, + {name: "app/", isDir: true, mode: 0755}, + {name: "app/main.txt", content: "v1", mode: 0644}, + {name: "data/", isDir: true, mode: 0755}, + {name: "data/old.txt", content: "stale", mode: 0644}, + {name: "replacedir/", isDir: true, mode: 0755}, + {name: "replacedir/inner.txt", content: "inner", mode: 0644}, + }) + top := specLayer(t, []tarEntrySpec{ + {name: "etc/.wh.config.txt", content: "", mode: 0644}, + {name: "app/main.txt", content: "v2", mode: 0644}, + {name: "data/.wh..wh..opq", content: "", mode: 0644}, + {name: "data/new.txt", content: "new", mode: 0644}, + {name: "bin/", isDir: true, mode: 0755}, + {name: "bin/tool", content: "tool", mode: 0755}, + {name: "replacedir", content: "now a file", mode: 0644}, + }) + + img, err := mutate.AppendLayers(empty.Image, base, top) + require.NoError(t, err) + return img +} + +// composeFixture composes the standard fixture image into the shared OCI cache +// and returns a client plus its validated manifest model. +func composeFixture(t *testing.T, p *paths.Paths) (*ociClient, string, *imageManifestModel) { + t.Helper() + + img := composeTestImage(t) + writeLayerTestLayout(t, p, img) + + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) + digest, err := img.Digest() + require.NoError(t, err) + tag := digestToLayoutTag(digest.String()) + bundle, err := client.extractOCIImageBundle(tag) + require.NoError(t, err) + return client, tag, bundle.Model +} + +func TestComposeRootfsWhiteoutsAndOrdering(t *testing.T) { + p := paths.New(t.TempDir()) + client, tag, model := composeFixture(t, p) + require.Len(t, model.Layers, 2) + + dest := filepath.Join(t.TempDir(), "rootfs") + require.NoError(t, client.composeRootfs(dest, tag, model)) + + // Whiteout removed the base entry. + _, err := os.Lstat(filepath.Join(dest, "etc", "config.txt")) + require.True(t, os.IsNotExist(err), "whiteout must delete the base entry") + + // Plain replacement. + data, err := os.ReadFile(filepath.Join(dest, "app", "main.txt")) + require.NoError(t, err) + require.Equal(t, "v2", string(data)) + + // Opaque directory masked the base content. + _, err = os.Lstat(filepath.Join(dest, "data", "old.txt")) + require.True(t, os.IsNotExist(err), "opaque marker must mask base contents") + data, err = os.ReadFile(filepath.Join(dest, "data", "new.txt")) + require.NoError(t, err) + require.Equal(t, "new", string(data)) + + // Directory replaced by a regular file. + info, err := os.Lstat(filepath.Join(dest, "replacedir")) + require.NoError(t, err) + require.False(t, info.IsDir()) + data, err = os.ReadFile(filepath.Join(dest, "replacedir")) + require.NoError(t, err) + require.Equal(t, "now a file", string(data)) + + // New entry present with its mode. + info, err = os.Stat(filepath.Join(dest, "bin", "tool")) + require.NoError(t, err) + require.Equal(t, os.FileMode(0755), info.Mode().Perm()) + + // No whiteout markers survive composition. + require.NoError(t, filepath.Walk(dest, func(path string, info os.FileInfo, err error) error { + require.NoError(t, err) + require.NotContains(t, info.Name(), whiteoutPrefix, "whiteout marker leaked into composed rootfs") + return nil + })) +} + +// zeroLayerModel returns a schema-valid manifest model with no layers. +func zeroLayerModel() *imageManifestModel { + return &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: "sha256:" + strings.Repeat("ab", 32), + Config: manifestConfigRef{Digest: "sha256:" + strings.Repeat("cd", 32)}, + Layers: make([]layerDescriptor, 0), + } +} + +func TestComposeRootfsEmptyLayers(t *testing.T) { + p := paths.New(t.TempDir()) + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) + model := zeroLayerModel() + dest := filepath.Join(t.TempDir(), "rootfs") + require.NoError(t, client.composeRootfs(dest, model.Digest, model)) + entries, err := os.ReadDir(dest) + require.NoError(t, err) + require.Empty(t, entries) +} + +func TestComposeRootfsInvalidModel(t *testing.T) { + p := paths.New(t.TempDir()) + client, tag, model := composeFixture(t, p) + + model.Config.DiffIDs = model.Config.DiffIDs[:1] + err := client.composeRootfs(filepath.Join(t.TempDir(), "rootfs"), tag, model) + require.ErrorContains(t, err, "1 diff ids for 2 layers") +} + +func TestComposeRootfsMissingBlob(t *testing.T) { + p := paths.New(t.TempDir()) + client, err := newOCIClient(p.SystemOCICache()) + require.NoError(t, err) + + digestHex := "sha256:" + strings.Repeat("ab", 32) + model := &imageManifestModel{ + SchemaVersion: manifestModelSchemaVersion, + Digest: digestHex, + Config: manifestConfigRef{ + Digest: "sha256:" + strings.Repeat("cd", 32), + DiffIDs: []string{"sha256:" + strings.Repeat("ef", 32)}, + }, + Layers: []layerDescriptor{{ + Digest: "sha256:" + strings.Repeat("01", 32), + MediaType: "application/vnd.oci.image.layer.v1.tar+gzip", + DiffID: "sha256:" + strings.Repeat("ef", 32), + }}, + } + err = client.composeRootfs(t.TempDir(), digestHex, model) + require.ErrorContains(t, err, "missing from oci cache") +} + +func TestComposeRootfsDiffIDMismatch(t *testing.T) { + p := paths.New(t.TempDir()) + client, tag, model := composeFixture(t, p) + + // Replace the top layer's cached blob with different content so the + // unpacked diff id no longer matches the descriptor. + other := specLayer(t, []tarEntrySpec{{name: "other.txt", content: "other", mode: 0644}}) + otherBlob, err := other.Compressed() + require.NoError(t, err) + data, err := io.ReadAll(otherBlob) + require.NoError(t, err) + topHex := strings.TrimPrefix(model.Layers[1].Digest, "sha256:") + require.NoError(t, os.WriteFile(p.OCICacheBlob(topHex), data, 0644)) + + err = client.composeRootfs(filepath.Join(t.TempDir(), "rootfs"), tag, model) + require.ErrorContains(t, err, "diff id mismatch") +} + +// TestComposeRootfsExportsValidErofs composes the fixture image and exports it +// to erofs, then verifies the filesystem is intact and its contents match the +// composed tree. +func TestComposeRootfsExportsValidErofs(t *testing.T) { + if _, err := exec.LookPath("mkfs.erofs"); err != nil { + t.Skip("mkfs.erofs not available") + } + if _, err := exec.LookPath("fsck.erofs"); err != nil { + t.Skip("fsck.erofs not available") + } + + p := paths.New(t.TempDir()) + client, tag, model := composeFixture(t, p) + + staging := filepath.Join(t.TempDir(), "rootfs") + require.NoError(t, client.composeRootfs(staging, tag, model)) + + diskPath := filepath.Join(t.TempDir(), "rootfs.erofs") + size, err := ExportRootfs(staging, diskPath, FormatErofs) + require.NoError(t, err) + require.Greater(t, size, int64(0)) + + output, err := exec.Command("fsck.erofs", "--extract", diskPath).CombinedOutput() + require.NoError(t, err, "fsck.erofs failed: %s", output) +} diff --git a/lib/images/disk_usage.go b/lib/images/disk_usage.go index b65dba1ac..0cb6a74d6 100644 --- a/lib/images/disk_usage.go +++ b/lib/images/disk_usage.go @@ -3,8 +3,10 @@ package images import ( "encoding/json" "fmt" + "io/fs" "os" "path/filepath" + "strings" "syscall" ) @@ -82,6 +84,38 @@ func totalReadyImageBytesFromMetadata(imagesDir string) (int64, error) { return total, nil } +// totalLayerArtifactBytesFromFilesystem sums materialized layer artifacts. +func totalLayerArtifactBytesFromFilesystem(layersDir string) (int64, error) { + var total int64 + err := filepath.WalkDir(layersDir, func(path string, d fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if d.IsDir() { + if strings.HasPrefix(d.Name(), ".") && path != layersDir { + return filepath.SkipDir + } + return nil + } + if !strings.HasPrefix(d.Name(), "layer.") { + return nil + } + info, err := d.Info() + if err != nil { + return err + } + total += info.Size() + return nil + }) + if err != nil { + return 0, fmt.Errorf("walk layer artifacts: %w", err) + } + return total, nil +} + // totalOCICacheBlobBytesFromFilesystem sums blob sizes directly from the OCI cache blob store. // This counts the actual bytes on disk, including any blob files that are currently // present but no longer referenced by the OCI layout index. @@ -158,7 +192,11 @@ func (m *manager) computeDiskUsageTotals() (int64, int64, error) { if err != nil { return 0, 0, err } - return readyImageBytes, ociCacheBytes, nil + layerArtifactBytes, err := totalLayerArtifactBytesFromFilesystem(m.paths.ImageLayersDir()) + if err != nil { + return 0, 0, err + } + return readyImageBytes, ociCacheBytes + layerArtifactBytes, nil } func totalRootfsBytesInDigestDir(digestDir string) (int64, error) { diff --git a/lib/images/disk_usage_test.go b/lib/images/disk_usage_test.go index bd6056fef..7de35871b 100644 --- a/lib/images/disk_usage_test.go +++ b/lib/images/disk_usage_test.go @@ -64,6 +64,26 @@ func TestTotalReadyImageBytesFromMetadata_DeduplicatesMalformedAliases(t *testin require.Equal(t, int64(len("shared-rootfs")), total) } +func TestTotalLayerArtifactBytesFromFilesystem(t *testing.T) { + t.Parallel() + + layersDir := t.TempDir() + digestDir := filepath.Join(layersDir, "a", "b") + require.NoError(t, os.MkdirAll(digestDir, 0o755)) + require.NoError(t, os.WriteFile(filepath.Join(digestDir, "layer.erofs"), []byte("erofs"), 0o644)) + require.NoError(t, os.WriteFile(filepath.Join(digestDir, "layer.ext4"), []byte("ext4"), 0o644)) + require.NoError(t, os.WriteFile(filepath.Join(digestDir, "artifact.erofs.json"), []byte("record"), 0o644)) + + // In-progress temp dirs should be skipped. + unpackDir := filepath.Join(digestDir, ".unpack-tmp") + require.NoError(t, os.MkdirAll(unpackDir, 0o755)) + require.NoError(t, os.WriteFile(filepath.Join(unpackDir, "layer.bin"), []byte("unpacked-content"), 0o644)) + + total, err := totalLayerArtifactBytesFromFilesystem(layersDir) + require.NoError(t, err) + require.Equal(t, int64(len("erofs")+len("ext4")), total) +} + func TestTotalReadyImageBytesFromMetadata_UsesRootfsFallbackForReadyImageWithoutSize(t *testing.T) { t.Parallel() diff --git a/lib/images/layer_artifact.go b/lib/images/layer_artifact.go new file mode 100644 index 000000000..0f40a0960 --- /dev/null +++ b/lib/images/layer_artifact.go @@ -0,0 +1,394 @@ +package images + +import ( + "compress/gzip" + "context" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "io" + "io/fs" + "os" + "path/filepath" + "strings" + "time" + + "github.com/kernel/hypeman/lib/paths" + "github.com/klauspost/compress/zstd" + "github.com/opencontainers/umoci/oci/layer" +) + +// whiteoutPrefix marks OCI whiteout entries (".wh." and ".wh..wh..opq"). +// umoci interprets them during extraction: DirRootfs applies them against the +// tree being composed, OverlayfsRootfs converts them to overlayfs whiteout +// inodes and opaque xattrs so per-layer artifacts can later be stacked. +const whiteoutPrefix = ".wh." + +const ( + layerRecordSchemaVersion = 1 + maxLayerUnpackedBytes = 100 << 30 +) + +var errCorruptLayerRecord = errors.New("corrupt layer record") + +// layerArtifact is the persisted record for one materialized layer artifact. +// The key is the compressed layer blob digest plus the artifact format, so the +// same layer can coexist in several materializations. +type layerArtifact struct { + SchemaVersion int `json:"schema_version"` + Digest string `json:"digest"` // compressed layer blob digest, sha256:... + DiffID string `json:"diff_id,omitempty"` + Format string `json:"format"` + SizeBytes int64 `json:"size_bytes"` // artifact bytes on disk + UnpackedBytes int64 `json:"unpacked_bytes"` // decompressed tar stream bytes + CreatedAt time.Time `json:"created_at"` +} + +// validate checks a record read back from disk. The format fully determines +// the artifact options (erofs is always lz4-compressed, ext4 uncompressed), +// so only the format is stored. +func (a *layerArtifact) validate() error { + if a.SchemaVersion != layerRecordSchemaVersion { + return fmt.Errorf("unsupported schema version: %d", a.SchemaVersion) + } + if a.Digest == "" { + return fmt.Errorf("missing digest") + } + if a.Format != string(FormatErofs) && a.Format != string(FormatExt4) { + return fmt.Errorf("invalid format: %s", a.Format) + } + if a.SizeBytes < 0 || a.UnpackedBytes < 0 { + return fmt.Errorf("invalid size") + } + return nil +} + +func (a *layerArtifact) matches(desc layerDescriptor) bool { + if a.Digest != desc.Digest || a.Format != layerArtifactFormat() { + return false + } + return desc.DiffID == "" || a.DiffID == desc.DiffID +} + +func layerArtifactFormat() string { + switch DefaultImageFormat { + case FormatErofs, FormatExt4: + return string(DefaultImageFormat) + default: + return "" + } +} + +func layerArtifactPath(p *paths.Paths, layerHex string) string { + return p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat()) +} + +func layerArtifactRecordPath(p *paths.Paths, layerHex string) string { + return p.ImageLayerRecordForFormat(layerHex, layerArtifactFormat()) +} + +// layerMapOptions preserves tar ownership when running as root. Otherwise +// umoci's rootless mode skips chown and stands in empty files for device nodes. +func layerMapOptions() layer.MapOptions { + return layer.MapOptions{Rootless: os.Geteuid() != 0} +} + +// layerArtifactOnDiskFormat is the extraction format for per-layer artifacts. +// Whiteouts become overlayfs whiteout inodes and opaque xattrs, the form an +// overlayfs mount of stacked layers understands, so the artifact retains the +// layer's deletions without a private marker format. +func layerArtifactOnDiskFormat() layer.OnDiskFormat { + return layer.OverlayfsRootfs{MapOptions: layerMapOptions()} +} + +// composeOnDiskFormat applies whiteouts against the tree being composed. +func composeOnDiskFormat() layer.OnDiskFormat { + return layer.DirRootfs{MapOptions: layerMapOptions()} +} + +// readLayerRecord loads the artifact record for a layer digest, if present. +// A missing record returns (nil, nil): the layer simply was never +// materialized. +func readLayerRecord(p *paths.Paths, layerHex string) (*layerArtifact, error) { + data, err := os.ReadFile(layerArtifactRecordPath(p, layerHex)) + if err != nil { + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("read layer record: %w", err) + } + var record layerArtifact + if err := json.Unmarshal(data, &record); err != nil { + return nil, fmt.Errorf("%w: unmarshal layer record: %v", errCorruptLayerRecord, err) + } + if err := record.validate(); err != nil { + return nil, fmt.Errorf("%w: invalid layer record: %v", errCorruptLayerRecord, err) + } + return &record, nil +} + +func discardLayerCache(p *paths.Paths, layerHex string) error { + for _, path := range []string{ + layerArtifactRecordPath(p, layerHex), + layerArtifactPath(p, layerHex), + } { + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + return err + } + } + return nil +} + +// materializeLayerArtifact ensures a layer has a materialized artifact keyed +// by its blob digest, building it from the shared OCI cache blob when absent. +// The layer is unpacked into an isolated temp directory, converted to the +// default image format, and installed atomically. Normal failures remove the +// temp directory; lifecycle reconciliation removes stale temp directories after +// an interrupted build. No production caller yet: pull integration and +// composition land in later changes. +// +// Concurrent callers share one build. The build itself is detached from the +// initiating caller's cancellation so one cancelled pull cannot fail every +// other pull waiting on the same layer; each caller still returns as soon as +// its own context is done. +func (m *manager) materializeLayerArtifact(ctx context.Context, desc layerDescriptor) (*layerArtifact, error) { + key := desc.Digest + "\x00" + layerArtifactFormat() + buildCtx := context.WithoutCancel(ctx) + result := m.layerFlights.DoChan(key, func() (any, error) { + return m.materializeLayerArtifactOnce(buildCtx, desc) + }) + select { + case <-ctx.Done(): + return nil, ctx.Err() + case shared := <-result: + if shared.Err != nil { + return nil, shared.Err + } + return shared.Val.(*layerArtifact), nil + } +} + +func (m *manager) materializeLayerArtifactOnce(ctx context.Context, desc layerDescriptor) (*layerArtifact, error) { + if layerArtifactFormat() == "" { + return nil, fmt.Errorf("unsupported layer artifact format: %s", DefaultImageFormat) + } + layerHex := strings.TrimPrefix(desc.Digest, "sha256:") + if err := paths.ValidatePathComponent(layerHex); err != nil { + return nil, fmt.Errorf("invalid layer digest %s: %w", desc.Digest, err) + } + + if record, err := readLayerRecord(m.paths, layerHex); err != nil { + if !errors.Is(err, errCorruptLayerRecord) { + return nil, err + } + if discardErr := discardLayerCache(m.paths, layerHex); discardErr != nil { + return nil, fmt.Errorf("discard corrupt layer cache: %w", discardErr) + } + } 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. + } + + layerDir := m.paths.ImageLayerDir(layerHex) + if err := os.MkdirAll(layerDir, 0755); err != nil { + return nil, fmt.Errorf("create layer directory: %w", err) + } + unpackDir, err := os.MkdirTemp(layerDir, ".unpack-*") + if err != nil { + return nil, fmt.Errorf("create unpack directory: %w", err) + } + defer removePath(unpackDir) + + stats, err := unpackCachedLayer(ctx, m.paths.SystemOCICache(), desc, unpackDir, layerArtifactOnDiskFormat()) + if err != nil { + return nil, err + } + return m.installLayerArtifact(desc, layerHex, unpackDir, stats) +} + +func (m *manager) installLayerArtifact(desc layerDescriptor, layerHex, unpackDir string, stats *unpackStats) (*layerArtifact, error) { + record := &layerArtifact{ + SchemaVersion: layerRecordSchemaVersion, + Digest: desc.Digest, + DiffID: stats.diffID, + Format: layerArtifactFormat(), + UnpackedBytes: stats.unpackedBytes, + CreatedAt: time.Now(), + } + + if err := installAtomically(layerArtifactPath(m.paths, layerHex), func(path string) error { + size, err := ExportRootfs(unpackDir, path, DefaultImageFormat) + if err != nil { + return err + } + record.SizeBytes = size + return nil + }); err != nil { + return nil, fmt.Errorf("install layer artifact %s: %w", desc.Digest, err) + } + + data, err := json.MarshalIndent(record, "", " ") + if err != nil { + return nil, fmt.Errorf("marshal layer record: %w", err) + } + if err := writeJSONAtomic(layerArtifactRecordPath(m.paths, layerHex), data); err != nil { + _ = os.Remove(layerArtifactPath(m.paths, layerHex)) + return nil, fmt.Errorf("write layer record: %w", err) + } + return record, nil +} + +type unpackStats struct { + unpackedBytes int64 + diffID string +} + +// contextReader fails reads once ctx is done so a cancelled caller stops a +// long extraction instead of running it to completion. +type contextReader struct { + ctx context.Context + reader io.Reader +} + +func (r contextReader) Read(p []byte) (int, error) { + if err := r.ctx.Err(); err != nil { + return 0, err + } + return r.reader.Read(p) +} + +// unpackCachedLayer locates desc's blob in the shared OCI cache, unpacks it +// into dest, and verifies the diff ID when the descriptor carries one. +func unpackCachedLayer(ctx context.Context, cacheDir string, desc layerDescriptor, dest string, onDisk layer.OnDiskFormat) (*unpackStats, error) { + layerHex := strings.TrimPrefix(desc.Digest, "sha256:") + if err := paths.ValidatePathComponent(layerHex); err != nil { + return nil, fmt.Errorf("invalid layer digest: %s", desc.Digest) + } + blobPath := filepath.Join(cacheDir, "blobs", "sha256", layerHex) + if _, err := os.Stat(blobPath); err != nil { + if os.IsNotExist(err) { + return nil, fmt.Errorf("layer blob missing from oci cache: %s", desc.Digest) + } + return nil, fmt.Errorf("stat layer blob: %w", err) + } + stats, err := unpackLayerBlob(ctx, blobPath, desc.MediaType, dest, onDisk) + if err != nil { + return nil, fmt.Errorf("unpack layer %s: %w", desc.Digest, err) + } + if desc.DiffID != "" && stats.diffID != desc.DiffID { + return nil, fmt.Errorf("layer %s diff id mismatch: got %s, want %s", desc.Digest, stats.diffID, desc.DiffID) + } + return stats, nil +} + +// unpackLayerBlob extracts one compressed layer blob into dest with umoci, +// which confines every entry to dest and interprets whiteouts per onDisk. +// The decompressed stream is hashed for the diff ID and capped in size. +func unpackLayerBlob(ctx context.Context, blobPath, mediaType, dest string, onDisk layer.OnDiskFormat) (*unpackStats, error) { + if err := os.MkdirAll(dest, 0755); err != nil { + return nil, fmt.Errorf("create extraction root: %w", err) + } + blob, err := os.Open(blobPath) + if err != nil { + return nil, fmt.Errorf("open blob: %w", err) + } + defer blob.Close() + + reader, closer, err := decompressLayer(blob, mediaType) + if err != nil { + return nil, err + } + defer closer.Close() + + hash := sha256.New() + limited := &io.LimitedReader{R: contextReader{ctx: ctx, reader: reader}, N: maxLayerUnpackedBytes + 1} + hashed := io.TeeReader(limited, hash) + if err := layer.UnpackLayer(dest, hashed, &layer.UnpackOptions{OnDiskFormat: onDisk}); err != nil { + return nil, err + } + if _, err := io.Copy(io.Discard, hashed); err != nil { + return nil, fmt.Errorf("drain layer: %w", err) + } + if limited.N == 0 { + return nil, fmt.Errorf("layer exceeds maximum unpacked size of %d bytes", maxLayerUnpackedBytes) + } + return &unpackStats{ + unpackedBytes: maxLayerUnpackedBytes + 1 - limited.N, + diffID: fmt.Sprintf("sha256:%x", hash.Sum(nil)), + }, nil +} + +// decompressLayer wraps the blob in the reader for its layer media type. Both +// OCI-style suffixes (+gzip, +zstd) and docker-style media types (tar.gzip, +// tar.zstd) are matched so neither encoding falls through to the raw path. +func decompressLayer(blob *os.File, mediaType string) (io.Reader, io.Closer, error) { + switch { + case strings.HasSuffix(mediaType, "+zstd"), strings.Contains(mediaType, "tar.zstd"): + decoder, err := zstd.NewReader(blob) + if err != nil { + return nil, nil, fmt.Errorf("zstd reader: %w", err) + } + return decoder, multiCloser{decoder.IOReadCloser(), blob}, nil + case strings.HasSuffix(mediaType, "+gzip"), strings.Contains(mediaType, "tar.gzip"): + gz, err := gzip.NewReader(blob) + if err != nil { + return nil, nil, fmt.Errorf("gzip reader: %w", err) + } + return gz, multiCloser{gz, blob}, nil + default: + return blob, blob, nil + } +} + +type multiCloser []io.Closer + +func (c multiCloser) Close() error { + var firstErr error + for _, closer := range c { + if err := closer.Close(); err != nil && firstErr == nil { + firstErr = err + } + } + return firstErr +} + +// removePath removes a tree that may contain read-only directories restored +// from layer metadata. +func removePath(path string) error { + if err := makeTreeWritable(path); err != nil && !os.IsNotExist(err) { + return err + } + if err := os.RemoveAll(path); err != nil && !os.IsNotExist(err) { + return err + } + return nil +} + +func makeTreeWritable(path string) error { + info, err := os.Lstat(path) + if err != nil { + return err + } + if !info.IsDir() || info.Mode()&os.ModeSymlink != 0 { + return nil + } + if err := os.Chmod(path, info.Mode().Perm()|0700); err != nil { + return err + } + return filepath.WalkDir(path, func(p string, entry fs.DirEntry, err error) error { + if err != nil { + return err + } + if entry.IsDir() && entry.Type()&os.ModeSymlink == 0 { + info, err := entry.Info() + if err != nil { + return err + } + return os.Chmod(p, info.Mode().Perm()|0700) + } + return nil + }) +} diff --git a/lib/images/layer_artifact_test.go b/lib/images/layer_artifact_test.go new file mode 100644 index 000000000..5e81a4624 --- /dev/null +++ b/lib/images/layer_artifact_test.go @@ -0,0 +1,405 @@ +package images + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "context" + "crypto/sha256" + "encoding/json" + "fmt" + "os" + "os/exec" + "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" + "golang.org/x/sys/unix" +) + +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) { + 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()), + }))) +} + +func layerDescFromImage(t *testing.T, img gcr.Image, index int) layerDescriptor { + t.Helper() + manifest, err := img.Manifest() + require.NoError(t, err) + configFile, err := img.ConfigFile() + require.NoError(t, err) + layer := manifest.Layers[index] + return layerDescriptor{ + Digest: layer.Digest.String(), + Size: layer.Size, + MediaType: string(layer.MediaType), + DiffID: configFile.RootFS.DiffIDs[index].String(), + } +} + +// tarGz builds a gzipped tar from the entries written by fn. +func tarGz(t *testing.T, fn func(tw *tar.Writer)) []byte { + t.Helper() + var buf bytes.Buffer + gzw := gzip.NewWriter(&buf) + tw := tar.NewWriter(gzw) + fn(tw) + require.NoError(t, tw.Close()) + require.NoError(t, gzw.Close()) + return buf.Bytes() +} + +func writeTarEntry(t *testing.T, tw *tar.Writer, header *tar.Header, content string) { + t.Helper() + if header.Typeflag == tar.TypeReg { + header.Size = int64(len(content)) + } + require.NoError(t, tw.WriteHeader(header)) + if content != "" { + _, err := tw.Write([]byte(content)) + require.NoError(t, err) + } +} + +func fileEntry(name, content string) func(*testing.T, *tar.Writer) { + return func(t *testing.T, tw *tar.Writer) { + writeTarEntry(t, tw, &tar.Header{Name: name, Typeflag: tar.TypeReg, Mode: 0644}, content) + } +} + +func dirEntry(name string) func(*testing.T, *tar.Writer) { + return func(t *testing.T, tw *tar.Writer) { + writeTarEntry(t, tw, &tar.Header{Name: name, Typeflag: tar.TypeDir, Mode: 0755}, "") + } +} + +func symlinkEntry(name, target string) func(*testing.T, *tar.Writer) { + return func(t *testing.T, tw *tar.Writer) { + writeTarEntry(t, tw, &tar.Header{Name: name, Typeflag: tar.TypeSymlink, Linkname: target}, "") + } +} + +func hardlinkEntry(name, target string) func(*testing.T, *tar.Writer) { + return func(t *testing.T, tw *tar.Writer) { + writeTarEntry(t, tw, &tar.Header{Name: name, Typeflag: tar.TypeLink, Linkname: target}, "") + } +} + +// writeLayerBlob writes a gzipped tar layer to a file and returns its path. +func writeLayerBlob(t *testing.T, dir, name string, entries ...func(*testing.T, *tar.Writer)) string { + t.Helper() + data := tarGz(t, func(tw *tar.Writer) { + for _, entry := range entries { + entry(t, tw) + } + }) + path := filepath.Join(dir, name) + require.NoError(t, os.WriteFile(path, data, 0644)) + return path +} + +// unpackInto applies a layer blob onto dest with compose semantics. +func unpackInto(t *testing.T, blobPath, dest string) *unpackStats { + t.Helper() + stats, err := unpackLayerBlob(context.Background(), blobPath, testTarGzMediaType, dest, composeOnDiskFormat()) + require.NoError(t, err) + return stats +} + +func requireNoWhiteoutMarkers(t *testing.T, root string) { + t.Helper() + require.NoError(t, filepath.WalkDir(root, func(path string, entry os.DirEntry, err error) error { + require.NoError(t, err) + require.NotContains(t, entry.Name(), whiteoutPrefix, "whiteout marker leaked into tree") + return nil + })) +} + +func TestMaterializeLayerArtifact(t *testing.T) { + if _, err := exec.LookPath("mkfs.erofs"); err != nil { + t.Skip("mkfs.erofs not available") + } + + p := paths.New(t.TempDir()) + img, err := mutate.AppendLayers(empty.Image, syntheticLayer(t, "base.txt", "base layer content")) + require.NoError(t, err) + writeLayerTestLayout(t, p, img) + + desc := layerDescFromImage(t, img, 0) + m := &manager{paths: p} + + record, err := m.materializeLayerArtifact(context.Background(), desc) + require.NoError(t, err) + require.Equal(t, desc.Digest, record.Digest) + require.Equal(t, desc.DiffID, record.DiffID) + require.Equal(t, layerArtifactFormat(), record.Format) + require.Greater(t, record.SizeBytes, int64(0)) + require.Greater(t, record.UnpackedBytes, int64(0)) + + layerHex := desc.Digest[len("sha256:"):] + artifactInfo, err := os.Stat(p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) + require.NoError(t, err, "layer.erofs must be installed") + + // A second materialization reuses the existing artifact. + reused, err := m.materializeLayerArtifact(context.Background(), desc) + require.NoError(t, err) + require.True(t, record.CreatedAt.Equal(reused.CreatedAt), "reuse must return the stored record") + artifactInfoAfter, err := os.Stat(p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) + require.NoError(t, err) + require.Equal(t, artifactInfo.ModTime(), artifactInfoAfter.ModTime(), "reuse must not rebuild") +} + +func TestMaterializeLayerArtifactRecoversCorruptRecord(t *testing.T) { + if _, err := exec.LookPath("mkfs.ext4"); err != nil { + t.Skip("mkfs.ext4 not available") + } + originalFormat := DefaultImageFormat + DefaultImageFormat = FormatExt4 + t.Cleanup(func() { DefaultImageFormat = originalFormat }) + + p := paths.New(t.TempDir()) + img, err := mutate.AppendLayers(empty.Image, syntheticLayer(t, "base.txt", "base layer content")) + require.NoError(t, err) + writeLayerTestLayout(t, p, img) + + desc := layerDescFromImage(t, img, 0) + m := &manager{paths: p} + first, err := m.materializeLayerArtifact(context.Background(), desc) + require.NoError(t, err) + + layerHex := desc.Digest[len("sha256:"):] + require.NoError(t, os.WriteFile( + p.ImageLayerRecordForFormat(layerHex, layerArtifactFormat()), + []byte("{not valid json"), + 0600, + )) + + second, err := m.materializeLayerArtifact(context.Background(), desc) + require.NoError(t, err) + require.NotEqual(t, first.CreatedAt, second.CreatedAt, "corrupt record must trigger a rebuild") + require.FileExists(t, p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) + + data, err := os.ReadFile(p.ImageLayerRecordForFormat(layerHex, layerArtifactFormat())) + require.NoError(t, err) + var record layerArtifact + require.NoError(t, json.Unmarshal(data, &record)) + require.NoError(t, record.validate()) +} + +// TestMaterializeLayerArtifactOutlivesCancelledCaller checks that a caller +// abandoning a shared build does not abort the build for everyone else: the +// artifact still lands even though the initiating context was cancelled. +func TestMaterializeLayerArtifactOutlivesCancelledCaller(t *testing.T) { + if _, err := exec.LookPath("mkfs.ext4"); err != nil { + t.Skip("mkfs.ext4 not available") + } + originalFormat := DefaultImageFormat + DefaultImageFormat = FormatExt4 + t.Cleanup(func() { DefaultImageFormat = originalFormat }) + + p := paths.New(t.TempDir()) + img, err := mutate.AppendLayers(empty.Image, syntheticLayer(t, "base.txt", "base layer content")) + require.NoError(t, err) + writeLayerTestLayout(t, p, img) + desc := layerDescFromImage(t, img, 0) + m := &manager{paths: p} + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + _, err = m.materializeLayerArtifact(ctx, desc) + require.ErrorIs(t, err, context.Canceled) + + layerHex := desc.Digest[len("sha256:"):] + require.Eventually(t, func() bool { + record, err := readLayerRecord(p, layerHex) + return err == nil && record != nil + }, 30*time.Second, 50*time.Millisecond, "detached build must still install the artifact") + require.FileExists(t, p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) +} + +func TestMaterializeLayerArtifactMissingBlob(t *testing.T) { + p := paths.New(t.TempDir()) + m := &manager{paths: p} + + _, err := m.materializeLayerArtifact(context.Background(), layerDescriptor{ + Digest: "sha256:abababababababababababababababababababababababababababababababab", + MediaType: testTarGzMediaType, + }) + require.ErrorContains(t, err, "missing from oci cache") +} + +// TestUnpackLayerBlobArtifactFormatKeepsWhiteouts checks that the artifact +// extraction format records deletions in overlayfs form: a 0:0 character +// device for a whiteout and the opaque xattr for an opaque directory. +func TestUnpackLayerBlobArtifactFormatKeepsWhiteouts(t *testing.T) { + if os.Geteuid() != 0 { + t.Skip("overlayfs whiteouts need mknod and trusted xattrs") + } + root := t.TempDir() + blob := writeLayerBlob(t, root, "layer.tar.gz", + fileEntry("keep.txt", "keep"), + dirEntry("gone/"), + fileEntry("gone/.wh.deleted.txt", ""), + dirEntry("opq/"), + fileEntry("opq/.wh..wh..opq", ""), + fileEntry("opq/fresh.txt", "fresh"), + ) + dest := filepath.Join(root, "dest") + _, err := unpackLayerBlob(context.Background(), blob, testTarGzMediaType, dest, layerArtifactOnDiskFormat()) + require.NoError(t, err) + + var stat unix.Stat_t + require.NoError(t, unix.Lstat(filepath.Join(dest, "gone", "deleted.txt"), &stat)) + require.Equal(t, uint32(unix.S_IFCHR), stat.Mode&unix.S_IFMT, "whiteout must be a character device") + require.Equal(t, uint64(0), uint64(stat.Rdev), "whiteout device must be 0:0") + + value := make([]byte, 8) + n, err := unix.Lgetxattr(filepath.Join(dest, "opq"), "trusted.overlay.opaque", value) + require.NoError(t, err) + require.Equal(t, "y", string(value[:n])) + requireNoWhiteoutMarkers(t, dest) +} + +func TestUnpackLayerBlobAppliesWhiteoutsAcrossLayers(t *testing.T) { + root := t.TempDir() + base := writeLayerBlob(t, root, "base.tar.gz", + dirEntry("etc/"), + fileEntry("etc/config.txt", "original"), + dirEntry("data/"), + fileEntry("data/old.txt", "stale"), + dirEntry("replacedir/"), + fileEntry("replacedir/inner.txt", "inner"), + fileEntry("added/foo", "old"), + ) + top := writeLayerBlob(t, root, "top.tar.gz", + fileEntry("etc/.wh.config.txt", ""), + fileEntry("data/.wh..wh..opq", ""), + fileEntry("data/new.txt", "new"), + fileEntry("replacedir", "now a file"), + fileEntry("added/.wh.foo", ""), + fileEntry("added/foo", "new"), + ) + dest := filepath.Join(root, "dest") + unpackInto(t, base, dest) + unpackInto(t, top, dest) + + _, err := os.Lstat(filepath.Join(dest, "etc", "config.txt")) + require.True(t, os.IsNotExist(err), "whiteout must delete the lower entry") + _, err = os.Lstat(filepath.Join(dest, "data", "old.txt")) + require.True(t, os.IsNotExist(err), "opaque marker must mask lower contents") + data, err := os.ReadFile(filepath.Join(dest, "data", "new.txt")) + require.NoError(t, err) + require.Equal(t, "new", string(data)) + info, err := os.Lstat(filepath.Join(dest, "replacedir")) + require.NoError(t, err) + require.False(t, info.IsDir(), "file must replace lower directory") + data, err = os.ReadFile(filepath.Join(dest, "added", "foo")) + require.NoError(t, err) + require.Equal(t, "new", string(data), "whiteout then recreate in one layer keeps the new file") + requireNoWhiteoutMarkers(t, dest) +} + +// TestUnpackLayerBlobReplacesDirectorySymlink checks the standard runtime +// behavior: a directory in an upper layer replaces a symlink from a lower one +// rather than writing through it. +func TestUnpackLayerBlobReplacesDirectorySymlink(t *testing.T) { + root := t.TempDir() + base := writeLayerBlob(t, root, "base.tar.gz", + dirEntry("usr/"), + dirEntry("usr/bin/"), + fileEntry("usr/bin/sh", "sh"), + symlinkEntry("bin", "usr/bin"), + ) + top := writeLayerBlob(t, root, "top.tar.gz", + dirEntry("bin/"), + fileEntry("bin/tool", "tool"), + ) + dest := filepath.Join(root, "dest") + unpackInto(t, base, dest) + unpackInto(t, top, dest) + + info, err := os.Lstat(filepath.Join(dest, "bin")) + require.NoError(t, err) + require.True(t, info.IsDir(), "upper directory must replace the lower symlink") + require.FileExists(t, filepath.Join(dest, "bin", "tool")) + require.FileExists(t, filepath.Join(dest, "usr", "bin", "sh")) + require.NoFileExists(t, filepath.Join(dest, "usr", "bin", "tool")) +} + +func TestUnpackLayerBlobConfinesSymlinkTraversal(t *testing.T) { + root := t.TempDir() + blob := writeLayerBlob(t, root, "layer.tar.gz", + symlinkEntry("link", "../../escape"), + fileEntry("link/pwned", "x"), + ) + dest := filepath.Join(root, "dest") + unpackInto(t, blob, dest) + + require.NoFileExists(t, filepath.Join(root, "escape", "pwned")) + require.FileExists(t, filepath.Join(dest, "escape", "pwned"), "escaping link must be resolved inside the root") +} + +func TestUnpackLayerBlobPreservesHardlinks(t *testing.T) { + root := t.TempDir() + blob := writeLayerBlob(t, root, "layer.tar.gz", + fileEntry("a", "shared"), + hardlinkEntry("b", "a"), + ) + dest := filepath.Join(root, "dest") + unpackInto(t, blob, dest) + + a, err := os.Stat(filepath.Join(dest, "a")) + require.NoError(t, err) + b, err := os.Stat(filepath.Join(dest, "b")) + require.NoError(t, err) + require.True(t, os.SameFile(a, b)) +} + +func TestUnpackLayerBlobContextHonorsCancellation(t *testing.T) { + root := t.TempDir() + blob := writeLayerBlob(t, root, "layer.tar.gz", fileEntry("file", "x")) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + _, err := unpackLayerBlob(ctx, blob, testTarGzMediaType, filepath.Join(root, "dest"), composeOnDiskFormat()) + require.ErrorIs(t, err, context.Canceled) +} + +func TestUnpackLayerBlobIncludesTrailingTarPaddingInDiffID(t *testing.T) { + root := t.TempDir() + var tarData bytes.Buffer + tw := tar.NewWriter(&tarData) + writeTarEntry(t, tw, &tar.Header{Name: "file", Typeflag: tar.TypeReg, Mode: 0644}, "x") + require.NoError(t, tw.Close()) + _, err := tarData.Write(make([]byte, 512)) + require.NoError(t, err) + + var compressed bytes.Buffer + gzw := gzip.NewWriter(&compressed) + _, err = gzw.Write(tarData.Bytes()) + require.NoError(t, err) + require.NoError(t, gzw.Close()) + blob := filepath.Join(root, "layer.tar.gz") + require.NoError(t, os.WriteFile(blob, compressed.Bytes(), 0644)) + + stats := unpackInto(t, blob, filepath.Join(root, "dest")) + want := sha256.Sum256(tarData.Bytes()) + require.Equal(t, "sha256:"+fmt.Sprintf("%x", want), stats.diffID) + require.Equal(t, int64(tarData.Len()), stats.unpackedBytes) +} diff --git a/lib/images/layer_gc.go b/lib/images/layer_gc.go new file mode 100644 index 000000000..6392f19e3 --- /dev/null +++ b/lib/images/layer_gc.go @@ -0,0 +1,216 @@ +package images + +import ( + "context" + "fmt" + "io/fs" + "log/slog" + "os" + "path/filepath" + "strings" + "time" +) + +// layerEvictionGracePeriod keeps freshly written layer artifacts and temp +// directories out of cleanup so recovery and eviction never race builds that +// are still writing them. +const layerEvictionGracePeriod = 10 * time.Minute + +// referencedLayerDigests returns the set of layer blob digests referenced by +// the manifest models of every image in the content layout, plus the digests +// currently referenced by in-flight builds. Layer artifacts in this set are +// protected from eviction. Unreadable manifest models are skipped with a +// warning so one corrupt record cannot disable eviction entirely. +func (m *manager) referencedLayerDigests() map[string]struct{} { + refs := m.inflightLayerRefSnapshot() + contentRoot := filepath.Join(m.paths.ImagesDir(), "content") + err := filepath.WalkDir(contentRoot, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if entry.IsDir() || entry.Name() != "manifest.json" { + return nil + } + digestHex := filepath.Base(filepath.Dir(path)) + model, readErr := readManifestModel(m.paths, digestHex) + if readErr != nil { + slog.Warn("skipping unreadable manifest model for layer eviction", "digest", digestHex, "error", readErr) + return nil + } + if model == nil { + return nil + } + for _, layer := range model.Layers { + refs[strings.TrimPrefix(layer.Digest, "sha256:")] = struct{}{} + } + return nil + }) + if err != nil && !os.IsNotExist(err) { + slog.Warn("failed to walk content manifests for layer eviction", "error", err) + } + return refs +} + +// inflightLayerRefSnapshot returns the layer digests currently retained by +// in-flight builds. +func (m *manager) inflightLayerRefSnapshot() map[string]struct{} { + m.layerRefMu.Lock() + defer m.layerRefMu.Unlock() + refs := make(map[string]struct{}, len(m.inflightLayerRefs)) + for digestHex := range m.inflightLayerRefs { + refs[digestHex] = struct{}{} + } + return refs +} + +// reconcileLayerStore evicts unreferenced layer artifacts and refreshes the +// cached disk usage totals so accounting reflects the removals. +func (m *manager) reconcileLayerStore() { + m.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 := m.referencedLayerDigests() + + layersDir := m.paths.ImageLayersDir() + entries, err := os.ReadDir(layersDir) + if err != nil { + if !os.IsNotExist(err) { + slog.Warn("layer eviction failed to list layer store", "error", err) + } + return + } + + cutoff := time.Now().Add(-m.layerEvictionGrace) + evicted := 0 + var evictedBytes int64 + for _, entry := range entries { + if !entry.IsDir() { + continue + } + digestHex := entry.Name() + if _, referenced := refs[digestHex]; referenced { + continue + } + size, removed := m.tryEvictLayerArtifact(digestHex, filepath.Join(layersDir, digestHex), cutoff) + if !removed { + continue + } + evicted++ + evictedBytes += size + } + if evicted > 0 { + slog.Info("evicted unreferenced layer artifacts", "count", evicted, "bytes", evictedBytes) + if m.metrics != nil { + m.metrics.layerArtifactsEvicted.Add(context.Background(), int64(evicted)) + } + } +} + +// tryEvictLayerArtifact removes one unreferenced layer artifact if it is still +// stale and no build is materializing it. Reconciliation holds createMu while +// scanning and eviction, so new builds cannot begin without being retained. +func (m *manager) tryEvictLayerArtifact(digestHex, dirPath string, cutoff time.Time) (int64, bool) { + + // The candidate was selected outside the lock; re-check that a build has + // not retained the digest and the artifact has not been rewritten since. + if _, referenced := m.inflightLayerRefSnapshot()[digestHex]; referenced { + return 0, false + } + info, statErr := os.Stat(dirPath) + if statErr != nil || info.ModTime().After(cutoff) { + return 0, false + } + size, err := dirSize(dirPath) + if err != nil { + slog.Warn("failed to measure layer artifact size", "digest", digestHex, "error", err) + } + if err := os.RemoveAll(dirPath); err != nil { + slog.Warn("failed to evict unreferenced layer artifact", "digest", digestHex, "error", err) + return 0, false + } + return size, true +} + +// cleanStaleImageTempDirs removes temp directories left behind by builds that +// were interrupted mid-install, mid-materialization, or mid-tag promotion. +// Only directories older than the grace period are removed so live builds are +// never disturbed. +func (m *manager) cleanStaleImageTempDirs() { + roots := []string{ + m.paths.ImageLayersDir(), + filepath.Join(m.paths.ImagesDir(), "content"), + } + cutoff := time.Now().Add(-m.layerEvictionGrace) + for _, root := range roots { + err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if !entry.IsDir() { + return nil + } + name := entry.Name() + if !strings.HasPrefix(name, ".unpack-") && !strings.HasPrefix(name, ".install-") && !strings.HasPrefix(name, ".tag-stage-") { + return nil + } + info, statErr := os.Stat(path) + if statErr == nil && info.ModTime().Before(cutoff) { + _ = os.RemoveAll(path) + } + return fs.SkipDir + }) + if err != nil && !os.IsNotExist(err) { + slog.Warn("failed to clean stale image temp dirs", "root", root, "error", err) + } + } +} + +// totalLayerArtifactBytes sums the bytes held by materialized layer +// artifacts, matching what diskutilization.Collect counts for the same store. +func totalLayerArtifactBytes(layersDir string) (int64, error) { + var total int64 + err := filepath.WalkDir(layersDir, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if entry.IsDir() { + return nil + } + if !strings.HasPrefix(entry.Name(), "layer.") { + return nil + } + info, statErr := entry.Info() + if statErr != nil { + return nil + } + total += info.Size() + return nil + }) + if err != nil && !os.IsNotExist(err) { + return 0, fmt.Errorf("walk layer artifacts: %w", err) + } + return total, nil +} diff --git a/lib/images/lifecycle_test.go b/lib/images/lifecycle_test.go new file mode 100644 index 000000000..1519c618c --- /dev/null +++ b/lib/images/lifecycle_test.go @@ -0,0 +1,215 @@ +package images + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + gcr "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/empty" + "github.com/google/go-containerregistry/pkg/v1/layout" + "github.com/google/go-containerregistry/pkg/v1/mutate" + "github.com/kernel/hypeman/lib/paths" + "github.com/stretchr/testify/require" +) + +// writeSharedLayout writes several images into one OCI layout cache, each +// annotated with its own digest tag. +func writeSharedLayout(t *testing.T, p *paths.Paths, imgs ...gcr.Image) []string { + t.Helper() + + layoutPath, err := layout.Write(p.SystemOCICache(), empty.Index) + require.NoError(t, err) + + digests := make([]string, 0, len(imgs)) + for _, img := range imgs { + digest, err := img.Digest() + require.NoError(t, err) + require.NoError(t, layoutPath.AppendImage(img, layout.WithAnnotations(map[string]string{ + "org.opencontainers.image.ref.name": digestToLayoutTag(digest.String()), + }))) + digests = append(digests, digest.String()) + } + return digests +} + +func layerHexes(t *testing.T, p *paths.Paths) map[string]struct{} { + t.Helper() + entries, err := os.ReadDir(p.ImageLayersDir()) + require.NoError(t, err) + hexes := make(map[string]struct{}) + for _, entry := range entries { + if entry.IsDir() { + hexes[entry.Name()] = struct{}{} + } + } + return hexes +} + +// TestSharedLayersMaterializeOnceAndEvictWithReferences is the end-to-end +// lifecycle: two images share a base layer, the shared artifact is created +// once, survives the deletion of one image, and is evicted only when its last +// reference is gone. +func TestSharedLayersMaterializeOnceAndEvictWithReferences(t *testing.T) { + if _, err := exec.LookPath("mkfs.erofs"); err != nil { + t.Skip("mkfs.erofs not available") + } + dataDir := t.TempDir() + p := paths.New(dataDir) + mgr, err := NewManager(p, 1, nil) + require.NoError(t, err) + m := mgr.(*manager) + m.layerEvictionGrace = 0 + + base := syntheticLayer(t, "base.txt", "shared base content") + topA := syntheticLayer(t, "a.txt", "app A payload") + topB := syntheticLayer(t, "b.txt", "app B payload") + + imgA, err := mutate.AppendLayers(empty.Image, base, topA) + require.NoError(t, err) + imgB, err := mutate.AppendLayers(empty.Image, base, topB) + require.NoError(t, err) + + digests := writeSharedLayout(t, p, imgA, imgB) + digestA, digestB := digests[0], digests[1] + + baseManifest, err := imgA.Manifest() + require.NoError(t, err) + baseHex := baseManifest.Layers[0].Digest.Hex + topAHex := baseManifest.Layers[1].Digest.Hex + topBManifest, err := imgB.Manifest() + require.NoError(t, err) + topBHex := topBManifest.Layers[1].Digest.Hex + + ctx := context.Background() + const repoA = "kernel.local/apps/app-a" + const repoB = "kernel.local/apps/app-b" + + eventsA := make(chan StatusEvent, 2) + m.subscribeToReady(digestToLayoutTag(digestA), eventsA) + defer m.unsubscribeFromReady(digestToLayoutTag(digestA), eventsA) + _, err = m.ImportLocalImage(ctx, repoA, "v1", digestA) + require.NoError(t, err) + select { + case event := <-eventsA: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatal("image A did not become ready") + } + + eventsB := make(chan StatusEvent, 2) + m.subscribeToReady(digestToLayoutTag(digestB), eventsB) + defer m.unsubscribeFromReady(digestToLayoutTag(digestB), eventsB) + _, err = m.ImportLocalImage(ctx, repoB, "v1", digestB) + require.NoError(t, err) + select { + case event := <-eventsB: + require.Equal(t, StatusReady, event.Status) + case <-time.After(30 * time.Second): + t.Fatal("image B did not become ready") + } + + // The shared base layer materialized exactly once, alongside the two tops. + hexes := layerHexes(t, p) + require.Len(t, hexes, 3) + require.Contains(t, hexes, baseHex) + require.Contains(t, hexes, topAHex) + require.Contains(t, hexes, topBHex) + + // Deleting image A evicts only its unique layer; the shared base survives. + require.NoError(t, m.DeleteImage(ctx, repoA+"@"+digestA)) + hexes = layerHexes(t, p) + require.Len(t, hexes, 2) + require.Contains(t, hexes, baseHex, "shared base must survive while referenced") + require.Contains(t, hexes, topBHex) + require.NotContains(t, hexes, topAHex) + + // Deleting image B removes the last references: everything is evicted. + require.NoError(t, m.DeleteImage(ctx, repoB+"@"+digestB)) + hexes = layerHexes(t, p) + require.Empty(t, hexes, "unreferenced layer artifacts must be evicted") +} + +func newLifecycleTestManager(p *paths.Paths) *manager { + return &manager{ + paths: p, + inflightPulls: make(map[string]*inflightImagePull), + inflightLayerRefs: make(map[string]int), + } +} + +func TestTotalImageBytesIncludesLayerArtifacts(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+cacheBytes, totalBytes) +} + +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, + Config: manifestConfigRef{ + Digest: "sha256:" + strings.Repeat("c", 64), + DiffIDs: []string{"sha256:" + referencedHex}, + }, + Layers: []layerDescriptor{{Digest: "sha256:" + referencedHex, DiffID: "sha256:" + referencedHex}}, + } + require.NoError(t, writeManifestModel(p, referencedHex, model)) + require.NoError(t, os.MkdirAll(p.ImageLayerDir(referencedHex), 0o755)) + require.NoError(t, os.WriteFile(p.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 := 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 8b905a417..6eeaedf36 100644 --- a/lib/images/manager.go +++ b/lib/images/manager.go @@ -19,6 +19,7 @@ import ( "github.com/kernel/hypeman/lib/queue" "github.com/kernel/hypeman/lib/tags" "go.opentelemetry.io/otel/metric" + "golang.org/x/sync/singleflight" ) var errStaleBuild = errors.New("stale image build") @@ -55,7 +56,7 @@ type Manager interface { // TotalImageBytes returns the total size of all ready images on disk. // Used by the resource manager for disk capacity tracking. TotalImageBytes(ctx context.Context) (int64, error) - // TotalOCICacheBytes returns the total size of the OCI layer cache. + // TotalOCICacheBytes returns the total size of the OCI and materialized layer caches. // Used by the resource manager for disk capacity tracking. TotalOCICacheBytes(ctx context.Context) (int64, error) // WaitForReady blocks until the image identified by name reaches a terminal @@ -75,6 +76,10 @@ type manager struct { ociClient *ociClient queue *queue.Queue createMu sync.Mutex + layerFlights singleflight.Group + layerRefMu sync.Mutex + 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 @@ -103,6 +108,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), @@ -484,6 +491,11 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials } m.recordPullMetrics(ctx, "success") + if err := m.materializeLayerArtifacts(ctx, result); err != nil { + m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("materialize layers: %w", err), buildID) + return + } + // 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 { @@ -531,6 +543,35 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials buildStatus = "success" } +func (m *manager) materializeLayerArtifacts(ctx context.Context, result *pullResult) error { + if result == nil || result.Manifest == nil { + return nil + } + for _, desc := range result.Manifest.Layers { + digestHex := strings.TrimPrefix(desc.Digest, "sha256:") + m.createMu.Lock() + m.layerRefMu.Lock() + m.inflightLayerRefs[digestHex]++ + m.layerRefMu.Unlock() + m.createMu.Unlock() + + _, err := m.materializeLayerArtifact(ctx, desc) + + m.createMu.Lock() + m.layerRefMu.Lock() + m.inflightLayerRefs[digestHex]-- + if m.inflightLayerRefs[digestHex] == 0 { + delete(m.inflightLayerRefs, digestHex) + } + m.layerRefMu.Unlock() + m.createMu.Unlock() + if err != nil { + return err + } + } + return nil +} + func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID, diskTempPath string) error { m.createMu.Lock() defer m.createMu.Unlock() @@ -847,14 +888,14 @@ func (m *manager) deleteTaggedImage(repository, tag 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, cacheBytes, err := m.getDiskUsageTotals() if err != nil { return 0, err } - return readyImageBytes, nil + return readyImageBytes + cacheBytes, nil } -// TotalOCICacheBytes returns the total size of the OCI layer cache. +// TotalOCICacheBytes returns the total size of the OCI and materialized layer caches. func (m *manager) TotalOCICacheBytes(ctx context.Context) (int64, error) { _, ociCacheBytes, err := m.getDiskUsageTotals() if err != nil { diff --git a/lib/images/metrics.go b/lib/images/metrics.go index d860885b2..75164a6cd 100644 --- a/lib/images/metrics.go +++ b/lib/images/metrics.go @@ -11,11 +11,12 @@ import ( // Metrics holds the metrics instruments for image operations. type Metrics struct { - buildDuration metric.Float64Histogram - buildPhaseDuration metric.Float64Histogram - ociLayerCount metric.Int64Histogram - ociCompressedBytes metric.Int64Histogram - pullsTotal metric.Int64Counter + buildDuration metric.Float64Histogram + buildPhaseDuration metric.Float64Histogram + ociLayerCount metric.Int64Histogram + ociCompressedBytes metric.Int64Histogram + pullsTotal metric.Int64Counter + layerArtifactsEvicted metric.Int64Counter } // newMetrics creates and registers all image metrics. @@ -78,6 +79,14 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) { return nil, err } + layerArtifactsEvicted, err := meter.Int64Counter( + "hypeman_images_layer_artifacts_evicted_total", + metric.WithDescription("Total number of shared layer artifacts evicted after their last reference was removed"), + ) + if err != nil { + return nil, err + } + // Register observable gauges for queue length and total images buildQueueLength, err := meter.Int64ObservableGauge( "hypeman_images_build_queue_length", @@ -123,11 +132,12 @@ func newMetrics(meter metric.Meter, m *manager) (*Metrics, error) { } return &Metrics{ - buildDuration: buildDuration, - buildPhaseDuration: buildPhaseDuration, - ociLayerCount: ociLayerCount, - ociCompressedBytes: ociCompressedBytes, - pullsTotal: pullsTotal, + buildDuration: buildDuration, + buildPhaseDuration: buildPhaseDuration, + ociLayerCount: ociLayerCount, + ociCompressedBytes: ociCompressedBytes, + pullsTotal: pullsTotal, + layerArtifactsEvicted: layerArtifactsEvicted, }, nil } diff --git a/lib/images/oci.go b/lib/images/oci.go index 88bd6de24..bb5ac888a 100644 --- a/lib/images/oci.go +++ b/lib/images/oci.go @@ -283,9 +283,9 @@ func (c *ociClient) pullAndExportWithPlatformAuth(ctx context.Context, imageRef, result.LayerCount = bundle.LayerCount result.CompressedBytes = bundle.CompressedBytes - // Unpack layers to the export directory + // Compose the rootfs from the shared layer blobs in manifest order. if err := result.measure("layer_unpack", func() error { - return c.unpackLayers(ctx, layoutTag, exportDir) + return c.composeRootfsContext(ctx, exportDir, layoutTag, bundle.Model) }); err != nil { return result, fmt.Errorf("unpack layers: %w", err) } 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 { diff --git a/lib/images/testlayers_test.go b/lib/images/testlayers_test.go new file mode 100644 index 000000000..d56cf1015 --- /dev/null +++ b/lib/images/testlayers_test.go @@ -0,0 +1,52 @@ +package images + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "io" + "testing" + + gcr "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/tarball" + "github.com/stretchr/testify/require" +) + +type tarEntrySpec struct { + name string + content string + isDir bool + mode int64 +} + +// specLayer builds a gzipped tar layer from entry specs in order. +func specLayer(t *testing.T, entries []tarEntrySpec) gcr.Layer { + t.Helper() + + var buf bytes.Buffer + gzw := gzip.NewWriter(&buf) + tw := tar.NewWriter(gzw) + for _, entry := range entries { + if entry.isDir { + require.NoError(t, tw.WriteHeader(&tar.Header{Name: entry.name, Typeflag: tar.TypeDir, Mode: entry.mode})) + continue + } + require.NoError(t, tw.WriteHeader(&tar.Header{ + Name: entry.name, + Typeflag: tar.TypeReg, + Mode: entry.mode, + Size: int64(len(entry.content)), + })) + _, err := tw.Write([]byte(entry.content)) + require.NoError(t, err) + } + require.NoError(t, tw.Close()) + require.NoError(t, gzw.Close()) + + data := buf.Bytes() + layer, err := tarball.LayerFromOpener(func() (io.ReadCloser, error) { + return io.NopCloser(bytes.NewReader(data)), nil + }) + require.NoError(t, err) + return layer +} diff --git a/lib/paths/paths.go b/lib/paths/paths.go index a6218cadd..67acbad9b 100644 --- a/lib/paths/paths.go +++ b/lib/paths/paths.go @@ -177,6 +177,28 @@ func (p *Paths) ImageRepositoryTagSymlink(repository, tag string) string { return filepath.Join(p.ImageRepositoriesDir(), repository, tag) } +// ImageLayersDir returns the root directory of the per-layer artifact store. +// Layer artifacts are content-addressed by the compressed layer blob digest. +func (p *Paths) ImageLayersDir() string { + return filepath.Join(p.dataDir, "images", "layers") +} + +// ImageLayerDir returns the artifact directory for one layer digest. +func (p *Paths) ImageLayerDir(layerHex string) string { + return filepath.Join(p.ImageLayersDir(), layerHex) +} + +// ImageLayerArtifactForFormat returns the path to a materialized layer artifact. +func (p *Paths) ImageLayerArtifactForFormat(layerHex, format string) string { + return filepath.Join(p.ImageLayerDir(layerHex), "layer."+format) +} + +// ImageLayerRecordForFormat returns the path to the artifact record for one +// materialized layer and format. +func (p *Paths) ImageLayerRecordForFormat(layerHex, format string) string { + return filepath.Join(p.ImageLayerDir(layerHex), "artifact."+format+".json") +} + // ImageDigestDir returns the directory for a specific image digest. func (p *Paths) ImageDigestDir(repository, digestHex string) string { return filepath.Join(p.dataDir, "images", repository, digestHex) diff --git a/lib/resources/resource.go b/lib/resources/resource.go index 86f644bda..1a3742bd8 100644 --- a/lib/resources/resource.go +++ b/lib/resources/resource.go @@ -132,7 +132,7 @@ type InstanceAllocation struct { type ImageLister interface { // TotalImageBytes returns the total size of all images on disk. TotalImageBytes(ctx context.Context) (int64, error) - // TotalOCICacheBytes returns the total size of the OCI layer cache. + // TotalOCICacheBytes returns the total size of the OCI and materialized layer caches. TotalOCICacheBytes(ctx context.Context) (int64, error) }