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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 36 additions & 1 deletion pkg/clip/clip.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"fmt"
"os"
"path/filepath"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -211,13 +212,22 @@ func MountArchive(options MountOptions) (func() error, <-chan error, *fuse.Serve
}

root, _ := clipfs.Root()
maxWrite := 1024 * 1024
if limit := spliceSafeMaxWrite(); limit > 0 && maxWrite > limit {
// go-fuse sets max_read = MaxWrite and serves fd-backed reads by
// splicing header+payload+one page through a pipe bounded by
// fs.pipe-max-size. If a full read doesn't fit, every read falls back
// to a copy and go-fuse logs "trySplice: splice: want N bytes".
log.Info().Int("max_write", limit).Msg("capping FUSE read/write size to fit fs.pipe-max-size")
maxWrite = limit
}
server, err := fuse.NewServer(fs.NewNodeFS(root, immutableFilesystemOptions()), options.MountPoint, &fuse.MountOptions{
MaxBackground: 512,
DisableXAttrs: true,
EnableSymlinkCaching: true,
SyncRead: false,
RememberInodes: true,
MaxWrite: 1024 * 1024,
MaxWrite: maxWrite,
MaxReadAhead: 1024 * 1024,
})
if err != nil {
Expand Down Expand Up @@ -345,6 +355,7 @@ type CreateFromOCIImageOptions struct {
Platform *v1.Platform
ContentCache storage.ContentCache // Optional cache to warm with decompressed layer streams
ContentCacheDir string // Optional temp directory for layer cache upload spooling
SeedDecompressed bool // Store freshly indexed layers' decompressed bytes in ContentCache
LayerIndexCache storage.LayerIndexCache // Optional per-layer index artifact cache (skips pull+index on hit)
IndexConcurrency int // Max layers indexed concurrently (default 4)
}
Expand Down Expand Up @@ -380,6 +391,7 @@ func CreateFromOCIImage(ctx context.Context, options CreateFromOCIImageOptions)
Platform: options.Platform,
ContentCache: options.ContentCache,
ContentCacheDir: options.ContentCacheDir,
SeedDecompressed: options.SeedDecompressed,
LayerIndexCache: options.LayerIndexCache,
IndexConcurrency: options.IndexConcurrency,
}, options.OutputPath)
Expand Down Expand Up @@ -428,3 +440,26 @@ func CreateAndUploadOCIArchive(ctx context.Context, options CreateFromOCIImageOp

return nil
}

// spliceSafeMaxWrite returns the largest page-aligned FUSE read/write size for
// which go-fuse's zero-copy splice path (header + payload + one extra page)
// fits in the kernel's maximum pipe size. Returns 0 if the limit is unknown.
func spliceSafeMaxWrite() int {
content, err := os.ReadFile("/proc/sys/fs/pipe-max-size")
if err != nil {
return 0
}
pipeMax, err := strconv.Atoi(strings.TrimSpace(string(content)))
if err != nil || pipeMax <= 0 {
return 0
}
return spliceSafeMaxWriteFor(pipeMax, os.Getpagesize())
}

func spliceSafeMaxWriteFor(pipeMax, pageSize int) int {
if pageSize <= 0 || pipeMax <= 2*pageSize {
return 0
}
limit := pipeMax - 2*pageSize
return limit - limit%pageSize
}
18 changes: 18 additions & 0 deletions pkg/clip/clip_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,3 +18,21 @@ func TestImmutableFilesystemOptionsCacheMetadata(t *testing.T) {
}
}
}

func TestSpliceSafeMaxWriteFor(t *testing.T) {
const page = 4096
for _, c := range []struct{ pipeMax, want int }{
{1 << 20, (1 << 20) - 2*page},
{4 << 20, (4 << 20) - 2*page},
{2 * page, 0},
{0, 0},
} {
got := spliceSafeMaxWriteFor(c.pipeMax, page)
if got != c.want {
t.Fatalf("spliceSafeMaxWriteFor(%d) = %d, want %d", c.pipeMax, got, c.want)
}
if got > 0 && 16+got+page > c.pipeMax {
t.Fatalf("pipeMax %d: %d leaves no room for header and extra page", c.pipeMax, got)
}
}
}
3 changes: 3 additions & 0 deletions pkg/clip/fsnode.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,9 @@ func (n *FSNode) readData(ctx context.Context, dest []byte, off int64) (fuse.Rea
nRead, err = n.filesystem.storage.ReadFile(n.clipNode, dest[:readLen], off)
}
if err != nil {
// The process sees EIO with no detail; leave the cause where an
// operator can find it.
log.Warn().Err(err).Str("path", n.clipNode.Path).Int64("offset", off).Int64("length", readLen).Msg("read failed, returning EIO to the container")
return nil, syscall.EIO
}
} else {
Expand Down
15 changes: 6 additions & 9 deletions pkg/clip/layer_artifact.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,15 +112,12 @@ func (ca *ClipArchiver) applyLayerArtifact(index *btree.BTree, artifact *LayerAr
case LayerEntryOpaqueWhiteout:
ca.deleteRange(index, entry.Path+"/")
case LayerEntryHardLink:
targetNode := index.Get(&common.ClipNode{Path: entry.Target})
if targetNode != nil {
tn := targetNode.(*common.ClipNode)
index.Set(&common.ClipNode{
Path: entry.Path,
NodeType: common.FileNode,
Attr: tn.Attr,
Remote: tn.Remote,
})
// The same inode under another name: a copy of the target, symlinks included
// (nix's optimised store hard-links symlinks into /nix/store/.links).
if targetNode := index.Get(&common.ClipNode{Path: entry.Target}); targetNode != nil {
linked := *targetNode.(*common.ClipNode)
linked.Path = entry.Path
index.Set(&linked)
}
}
}
Expand Down
19 changes: 19 additions & 0 deletions pkg/clip/layer_artifact_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,25 @@ func TestLayerArtifactRoundTripDeterminism(t *testing.T) {
assert.Equal(t, common.SymLinkNode, nodes["/link"].NodeType)
}

func TestLayerArtifactHardLinkToSymlink(t *testing.T) {
archiver := NewClipArchiver()

layer := buildLayer(t, []tarEntry{
{name: "dir/", typeflag: tar.TypeDir},
{name: "dir/a.txt", typeflag: tar.TypeReg, content: "hello"},
{name: "link", typeflag: tar.TypeSymlink, linkname: "dir/a.txt"},
{name: "hard-to-link", typeflag: tar.TypeLink, linkname: "link"},
})

index := archiver.newIndex()
archiver.applyLayerArtifact(index, indexLayerHelper(t, archiver, layer, "sha256:layer1"))

nodes := indexPaths(index)
require.Contains(t, nodes, "/hard-to-link")
assert.Equal(t, common.SymLinkNode, nodes["/hard-to-link"].NodeType, "a hard link to a symlink is a symlink")
assert.Equal(t, nodes["/link"].Target, nodes["/hard-to-link"].Target)
}

func TestLayerArtifactSanitizesUnsetTarTimes(t *testing.T) {
archiver := NewClipArchiver()

Expand Down
32 changes: 32 additions & 0 deletions pkg/clip/layer_blob_cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package clip
import (
"archive/tar"
"bytes"
"compress/gzip"
"context"
"crypto/sha256"
"encoding/hex"
Expand Down Expand Up @@ -190,3 +191,34 @@ func TestCompressedLayerContentCacheReadThrough(t *testing.T) {
assert.Equal(t, bytes1, bytes2, "index must be identical regardless of layer source")
assert.Equal(t, hashes1, hashes2)
}

func TestSeedDecompressedStoresIndexedLayerBytes(t *testing.T) {
compressed := buildLayer(t, []tarEntry{
{name: "dir/", typeflag: tar.TypeDir},
{name: "dir/a.txt", typeflag: tar.TypeReg, content: "hello"},
})
sum := sha256.Sum256(compressed)
digest := "sha256:" + hex.EncodeToString(sum[:])

for _, seed := range []bool{false, true} {
cache := newFakeBlobContentCache()
artifact, err := NewClipArchiver().indexLayerToArtifact(
context.Background(),
io.NopCloser(bytes.NewReader(compressed)),
digest,
IndexOCIImageOptions{CheckpointMiB: 2, ContentCache: cache, ContentCacheDir: t.TempDir(), SeedDecompressed: seed},
nil,
)
require.NoError(t, err)
assert.Equal(t, seed, cache.has(artifact.DecompressedHash), "seed=%v", seed)
if seed {
data, err := cache.GetContent(artifact.DecompressedHash, 0, artifact.UncompressedSize, struct{ RoutingKey string }{})
require.NoError(t, err)
gz, err := gzip.NewReader(bytes.NewReader(compressed))
require.NoError(t, err)
want, err := io.ReadAll(gz)
require.NoError(t, err)
assert.Equal(t, want, data)
}
}
}
88 changes: 84 additions & 4 deletions pkg/clip/oci_indexer.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"path"
"runtime"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"
Expand Down Expand Up @@ -71,8 +72,13 @@ type IndexOCIImageOptions struct {
Platform *v1.Platform // Target platform (defaults to linux/runtime.GOARCH)
ContentCache storage.ContentCache // optional remote cache for fully decompressed layers
ContentCacheDir string // optional temp directory for cache upload spooling
SeedDecompressed bool // store each freshly indexed layer's decompressed bytes in ContentCache
LayerIndexCache storage.LayerIndexCache // optional cache of per-layer index artifacts (skips pull+index on hit)
IndexConcurrency int // max layers indexed concurrently (default 4)

// seeder runs SeedDecompressed uploads off the indexing goroutines so a
// slow store does not hold an indexing slot; set by IndexOCIImage.
seeder *layerSeeder
}

const defaultIndexConcurrency = 4
Expand Down Expand Up @@ -384,6 +390,14 @@ func (ca *ClipArchiver) IndexOCIImage(ctx context.Context, opts IndexOCIImageOpt
g, gctx := errgroup.WithContext(ctx)
g.SetLimit(concurrency)

// Content cache seeds run alongside indexing, on the parent context
// (gctx is cancelled once g.Wait returns), and are all awaited before
// the index is returned: the indexing process may exit right after.
if opts.SeedDecompressed && opts.ContentCache != nil {
opts.seeder = newLayerSeeder(ctx, concurrency)
defer opts.seeder.wait()
}

var completedLayers atomic.Int64

for i := range layers {
Expand Down Expand Up @@ -557,12 +571,20 @@ func (ca *ClipArchiver) indexLayerToArtifact(
}
defer gzr.Close()

// Streaming hash computation via TeeReader.
// The runtime warms decompressed layers after first access; the build path
// keeps indexing strictly streaming so large layers do not pay extra disk
// writes just to seed an optional warm cache.
// Streaming hash computation via TeeReader. With SeedDecompressed the same
// decompressed stream is spooled to disk once and stored in the content
// cache after indexing, so the first container to use the layer reads it
// page-wise from the cache instead of materializing the whole layer.
hasher := sha256.New()
hashWriter := io.Writer(hasher)
var cacheSpool *indexedLayerContentCacheSpool
if opts.SeedDecompressed && opts.ContentCache != nil {
cacheSpool = newIndexedLayerContentCacheSpool(opts.ContentCacheDir, layerDigest)
if cacheSpool != nil {
defer cacheSpool.closeAndRemove()
hashWriter = io.MultiWriter(hasher, cacheSpool)
}
}
hashingReader := io.TeeReader(gzr, hashWriter)
uncompressedCounter := &countingReader{r: hashingReader, onRead: func(total int64) {
if onBytes != nil {
Expand Down Expand Up @@ -645,6 +667,36 @@ func (ca *ClipArchiver) indexLayerToArtifact(
// Finalize hash (includes all bytes: file contents + tar headers + padding)
decompressedHash := hex.EncodeToString(hasher.Sum(nil))

// The seed is awaited before indexing returns (the indexing process, a
// build worker, often exits right after); through the seeder it runs
// off this goroutine so the upload overlaps with indexing other layers.
if cacheSpool != nil {
if cacheSpool.err != nil {
log.Warn().Err(cacheSpool.err).Str("layer_digest", layerDigest).Msg("indexed layer not seeded: spool write failed")
} else if path, ok := cacheSpool.detach(); ok {
uncompressedBytes := uncompressedCounter.n
seed := func(ctx context.Context) {
seedStart := time.Now()
err := ca.storeIndexedLayerInContentCache(ctx, opts.ContentCache, path, decompressedHash, layerDigest)
os.Remove(path)
if err != nil {
log.Warn().Err(err).Str("layer_digest", layerDigest).Msg("indexed layer not seeded in content cache")
return
}
log.Info().
Str("layer_digest", layerDigest).
Int64("bytes", uncompressedBytes).
Dur("duration", time.Since(seedStart)).
Msg("seeded decompressed layer into content cache")
}
if opts.seeder != nil {
opts.seeder.run(seed)
} else {
seed(ctx)
}
}
}

return &LayerArtifact{
Version: LayerArtifactVersion,
LayerDigest: layerDigest,
Expand All @@ -656,6 +708,34 @@ func (ca *ClipArchiver) indexLayerToArtifact(
}, nil
}

// layerSeeder runs content cache seeds concurrently, up to a limit, and lets
// the indexer wait for all of them.
type layerSeeder struct {
ctx context.Context
wg sync.WaitGroup
sem chan struct{}
}

func newLayerSeeder(ctx context.Context, limit int) *layerSeeder {
if limit < 1 {
limit = 1
}
return &layerSeeder{ctx: ctx, sem: make(chan struct{}, limit)}
}

// run starts fn once a slot is free; it blocks only while every slot is busy.
func (s *layerSeeder) run(fn func(ctx context.Context)) {
s.wg.Add(1)
s.sem <- struct{}{}
go func() {
defer s.wg.Done()
defer func() { <-s.sem }()
fn(s.ctx)
}()
}

func (s *layerSeeder) wait() { s.wg.Wait() }

func (ca *ClipArchiver) storeIndexedLayerInContentCache(ctx context.Context, contentCache storage.ContentCache, filePath, decompressedHash, layerDigest string) error {
return ca.storeLayerBlobInContentCache(ctx, contentCache, filePath, decompressedHash, layerDigest, "indexed layer")
}
Expand Down
Loading
Loading