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..d7b79fe19 --- /dev/null +++ b/lib/images/layer_artifact.go @@ -0,0 +1,400 @@ +package images + +import ( + "compress/gzip" + "context" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "io" + "io/fs" + "log/slog" + "os" + "path/filepath" + "strings" + "time" + + "github.com/kernel/hypeman/lib/paths" + "github.com/klauspost/compress/zstd" + "github.com/opencontainers/umoci/oci/layer" +) + +const ( + layerRecordSchemaVersion = 1 + maxLayerUnpackedBytes = 100 << 30 + + // layerBuildTimeout bounds a shared layer build end to end. Builds are + // detached from the initiating request's context, so this deadline is the + // only thing that can abort a hung unpack and free the singleflight key. + layerBuildTimeout = time.Hour +) + +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 +} + +// matches decides whether the stored record satisfies a lookup. A descriptor +// without a DiffID (a manifest lookup that did not consult the image config) +// matches any record: the compressed digest content-addresses the blob, and +// unpackCachedLayer re-verifies the diff ID whenever one is supplied. +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. +// Unlike unpackLayers in oci.go, which maps container root to the current +// user, this deliberately leaves ownership untouched as root: artifacts must +// keep the layer's on-disk ownership for later stacking. +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()} +} + +// 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; 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. +// +// 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() + // The build outlives the initiating request, so the deadline below is its + // only bound: without it a hung cache-blob read would wedge the + // singleflight key, and every future caller for the layer, forever. + buildCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), layerBuildTimeout) + result := m.layerFlights.DoChan(key, func() (any, error) { + // Cancel here, in the flight's frame: the initiating caller may return + // long before the detached build finishes. + defer cancel() + 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 func() { + if err := removePath(unpackDir); err != nil { + slog.Warn("failed to remove layer unpack directory", "dir", unpackDir, "error", err) + } + }() + + stats, err := unpackCachedLayer(ctx, m.paths, 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 + blobDigest string // sha256 of the compressed bytes as read +} + +// 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 both the blob digest and the diff ID when the +// descriptor carries one. The caller must have validated desc.Digest. +func unpackCachedLayer(ctx context.Context, p *paths.Paths, desc layerDescriptor, dest string, onDisk layer.OnDiskFormat) (*unpackStats, error) { + blobPath := p.OCICacheBlob(strings.TrimPrefix(desc.Digest, "sha256:")) + 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.Digest != "" && stats.blobDigest != desc.Digest { + return nil, fmt.Errorf("layer blob digest mismatch: got %s, want %s", stats.blobDigest, desc.Digest) + } + 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 compressed stream is hashed to verify the blob digest, the decompressed +// stream 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() + + blobHash := sha256.New() + reader, closer, err := decompressLayer(io.TeeReader(blob, blobHash), 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)), + blobDigest: fmt.Sprintf("sha256:%x", blobHash.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(r io.Reader, mediaType string) (io.Reader, io.Closer, error) { + switch { + case strings.HasSuffix(mediaType, "+zstd"), strings.Contains(mediaType, "tar.zstd"): + decoder, err := zstd.NewReader(r) + if err != nil { + return nil, nil, fmt.Errorf("zstd reader: %w", err) + } + decodeCloser := decoder.IOReadCloser() + return decodeCloser, decodeCloser, nil + case strings.HasSuffix(mediaType, "+gzip"), strings.Contains(mediaType, "tar.gzip"): + gz, err := gzip.NewReader(r) + if err != nil { + return nil, nil, fmt.Errorf("gzip reader: %w", err) + } + return gz, gz, nil + default: + return r, io.NopCloser(r), nil + } +} + +// 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_smoke_test.go b/lib/images/layer_artifact_smoke_test.go new file mode 100644 index 000000000..ebbee631a --- /dev/null +++ b/lib/images/layer_artifact_smoke_test.go @@ -0,0 +1,210 @@ +package images + +import ( + "archive/tar" + "bytes" + "context" + "os" + "os/exec" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/google/go-containerregistry/pkg/v1/empty" + "github.com/google/go-containerregistry/pkg/v1/mutate" + "github.com/kernel/hypeman/lib/paths" + "github.com/klauspost/compress/zstd" + "github.com/stretchr/testify/require" + "golang.org/x/sys/unix" +) + +// Smoke tests for PR #456 edge cases beyond the standard suite. +// Media types, path validation, cache-miss rebuilds, singleflight sharing, +// device/fifo entries (root only), and disk-usage accounting. + +func writeRawTarLayer(t *testing.T, dir, name string, entries ...func(*testing.T, *tar.Writer)) string { + t.Helper() + var buf bytes.Buffer + tw := tar.NewWriter(&buf) + for _, entry := range entries { + entry(t, tw) + } + require.NoError(t, tw.Close()) + path := filepath.Join(dir, name) + require.NoError(t, os.WriteFile(path, buf.Bytes(), 0644)) + return path +} + +func writeZstdTarLayer(t *testing.T, dir, name string, entries ...func(*testing.T, *tar.Writer)) string { + t.Helper() + var tarBuf bytes.Buffer + tw := tar.NewWriter(&tarBuf) + for _, entry := range entries { + entry(t, tw) + } + require.NoError(t, tw.Close()) + + var zstBuf bytes.Buffer + zw, err := zstd.NewWriter(&zstBuf) + require.NoError(t, err) + _, err = zw.Write(tarBuf.Bytes()) + require.NoError(t, err) + require.NoError(t, zw.Close()) + path := filepath.Join(dir, name) + require.NoError(t, os.WriteFile(path, zstBuf.Bytes(), 0644)) + return path +} + +const ( + ociZstdMediaType = "application/vnd.oci.image.layer.v1.tar+zstd" + dockerZstdMediaType = "application/vnd.docker.image.rootfs.diff.tar.zstd" + dockerGzipMediaType = "application/vnd.docker.image.rootfs.diff.tar.gzip" + ociPlainTarMediaType = "application/vnd.oci.image.layer.v1.tar" +) + +func TestSmokeUnpackZstdMediaTypes(t *testing.T) { + root := t.TempDir() + for name, tc := range map[string]struct { + mediaType string + blob string + }{ + "oci_suffix_zstd": {ociZstdMediaType, writeZstdTarLayer(t, root, "oci.tar.zst", fileEntry("f.txt", "zstd"))}, + "docker_tar_zstd": {dockerZstdMediaType, writeZstdTarLayer(t, root, "docker.tar.zstd", fileEntry("f.txt", "zstd"))}, + "docker_tar_gzip": {dockerGzipMediaType, writeLayerBlob(t, root, "docker.tar.gz", fileEntry("f.txt", "gzip"))}, + "plain_tar": {ociPlainTarMediaType, writeRawTarLayer(t, root, "plain.tar", fileEntry("f.txt", "raw"))}, + "unknown_media_raw": {"application/vnd.custom.tar", writeRawTarLayer(t, root, "custom.tar", fileEntry("f.txt", "raw"))}, + } { + t.Run(name, func(t *testing.T) { + dest := filepath.Join(root, t.Name()) + stats, err := unpackLayerBlob(context.Background(), tc.blob, tc.mediaType, dest, composeOnDiskFormat()) + require.NoError(t, err) + require.Greater(t, stats.unpackedBytes, int64(0)) + require.FileExists(t, filepath.Join(dest, "f.txt")) + }) + } +} + +func TestSmokeRejectsBadLayerHex(t *testing.T) { + p := paths.New(t.TempDir()) + m := &manager{paths: p} + for _, digest := range []string{ + "sha256:..", + "sha256:../escape", + "sha256:a/b", + "sha256:", + } { + _, err := m.materializeLayerArtifact(context.Background(), layerDescriptor{ + Digest: digest, + MediaType: testTarGzMediaType, + }) + require.ErrorContains(t, err, "invalid layer digest", "digest %q must be rejected", digest) + } + require.NoDirExists(t, filepath.Join(p.ImageLayersDir(), "..", "layers-escape")) + + // A non-digest string without the sha256: prefix passes path validation + // and fails at blob lookup instead — acceptable, since real descriptors + // always carry manifest-parsed digests. + _, err := m.materializeLayerArtifact(context.Background(), layerDescriptor{ + Digest: "notadigest", + MediaType: testTarGzMediaType, + }) + require.ErrorContains(t, err, "missing from oci cache") +} + +func TestSmokeRebuildWhenArtifactFileMissing(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, "f.txt", "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:"):] + + // Record survives but the artifact file is gone: must rebuild, not reuse. + require.NoError(t, os.Remove(p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat()))) + second, err := m.materializeLayerArtifact(context.Background(), desc) + require.NoError(t, err) + require.True(t, second.CreatedAt.After(first.CreatedAt), "missing artifact must trigger a rebuild") + require.FileExists(t, p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) +} + +func TestSmokeConcurrentMaterializeSharesFlight(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, "f.txt", "shared")) + require.NoError(t, err) + writeLayerTestLayout(t, p, img) + desc := layerDescFromImage(t, img, 0) + m := &manager{paths: p} + + const callers = 8 + var wg sync.WaitGroup + records := make([]*layerArtifact, callers) + errs := make([]error, callers) + start := time.Now() + for i := 0; i < callers; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + records[i], errs[i] = m.materializeLayerArtifact(context.Background(), desc) + }(i) + } + wg.Wait() + require.NotEmpty(t, errs) + for i := range errs { + require.NoError(t, errs[i]) + require.Equal(t, records[0], records[i], "all callers must share one build result") + } + t.Logf("8 callers finished in %s sharing one flight", time.Since(start)) + layerHex := desc.Digest[len("sha256:"):] + require.FileExists(t, p.ImageLayerArtifactForFormat(layerHex, layerArtifactFormat())) +} + +func TestSmokeDeviceAndFifoEntries(t *testing.T) { + if os.Geteuid() != 0 { + t.Skip("device and fifo tar entries need root to mknod") + } + root := t.TempDir() + blob := writeRawTarLayer(t, root, "devices.tar", + fileEntry("regular.txt", "x"), + func(t *testing.T, tw *tar.Writer) { + require.NoError(t, tw.WriteHeader(&tar.Header{ + Name: "dev/null", Typeflag: tar.TypeChar, + Devmajor: 1, Devminor: 3, Mode: 0666, + })) + }, + func(t *testing.T, tw *tar.Writer) { + require.NoError(t, tw.WriteHeader(&tar.Header{ + Name: "pipes/fifo", Typeflag: tar.TypeFifo, Mode: 0644, + })) + }, + ) + dest := filepath.Join(root, "dest") + _, err := unpackLayerBlob(context.Background(), blob, ociPlainTarMediaType, dest, composeOnDiskFormat()) + require.NoError(t, err) + + require.FileExists(t, filepath.Join(dest, "regular.txt")) + var stat unix.Stat_t + require.NoError(t, unix.Lstat(filepath.Join(dest, "dev", "null"), &stat)) + require.Equal(t, uint32(unix.S_IFCHR), stat.Mode&unix.S_IFMT) + require.Equal(t, uint64(0x103), uint64(stat.Rdev), "char device 1:3") + require.NoError(t, unix.Lstat(filepath.Join(dest, "pipes", "fifo"), &stat)) + require.Equal(t, uint32(unix.S_IFIFO), stat.Mode&unix.S_IFMT) +} diff --git a/lib/images/layer_artifact_test.go b/lib/images/layer_artifact_test.go new file mode 100644 index 000000000..2ba416031 --- /dev/null +++ b/lib/images/layer_artifact_test.go @@ -0,0 +1,422 @@ +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/opencontainers/umoci/oci/layer" + "github.com/stretchr/testify/require" + "golang.org/x/sys/unix" +) + +// whiteoutPrefix marks OCI whiteout entries (".wh." and ".wh..wh..opq"). +// umoci interprets them during extraction: composeOnDiskFormat applies them +// against the tree being composed, layerArtifactOnDiskFormat converts them to +// overlayfs whiteout inodes and opaque xattrs so per-layer artifacts can +// later be stacked. +const whiteoutPrefix = ".wh." + +// composeOnDiskFormat applies whiteouts against the tree being composed. It +// belongs to the composition flow and moves to production with that change. +func composeOnDiskFormat() layer.OnDiskFormat { + return layer.DirRootfs{MapOptions: layerMapOptions()} +} + +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() + // The build usually loses the race against the cancelled context, but a + // fast build can finish first; either outcome is valid when detached. + if _, buildErr := m.materializeLayerArtifact(ctx, desc); buildErr != nil { + require.ErrorIs(t, buildErr, 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/manager.go b/lib/images/manager.go index 8b905a417..bceeb0f05 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,7 @@ type manager struct { ociClient *ociClient queue *queue.Queue createMu sync.Mutex + layerFlights singleflight.Group diskUsageMu sync.RWMutex tagGenerations map[string]uint64 requestedTags map[string]string // newest pull's digest per requested tag @@ -854,7 +856,7 @@ func (m *manager) TotalImageBytes(ctx context.Context) (int64, error) { return readyImageBytes, 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/paths/paths.go b/lib/paths/paths.go index fc3f221eb..b294f2417 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 d674c5d01..5943cb0bc 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) }