Skip to content

Commit a126feb

Browse files
committed
[CELL-297] nix-cache push uploads all chunks in parallel — eliminates per-chunk manifest churn that serialized 29 blob uploads on ultimate images
- feat(nixstore): chunked push splits tar to temp files then uploads all layers via single remote.Write with 4 parallel blob streams — ultimate image push no longer serializes 29 manifest writes + 29 re-fetches, cutting wall-clock time significantly - feat(nixstore): stall detector recognizes parallel upload finalization via TotalSizeHint byte comparison — no false stall abort while registry commits the single manifest after all blobs land - feat(nixstore): UploadJobs package var (default 4) controls parallel blob concurrency — tunable without code changes - docs(cli): nix-store push help text updated to describe parallel upload behavior — no functional change
1 parent 927ee0c commit a126feb

2 files changed

Lines changed: 84 additions & 41 deletions

File tree

‎cmd/nix_store.go‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,9 +47,10 @@ OCI tar+gzip layers atop --base, tagged --image, to the destination
4747
registry.
4848
4949
By default, the tar is split on entry boundaries into ~512 MB chunks,
50-
each uploaded as a separate layer. This reduces per-layer upload size,
51-
making GHCR throttling less likely and enabling partial retry on
52-
failure (CELL-297). Set DEVCELL_NIX_CHUNK_SIZE=0 to disable chunking.
50+
buffered to temp files, then uploaded in parallel (4 concurrent blob
51+
streams, 1 manifest write). This eliminates per-chunk manifest churn
52+
and overlaps gzip+upload across layers (CELL-297). Set
53+
DEVCELL_NIX_CHUNK_SIZE=0 to disable chunking.
5354
5455
Each layer is SINGLE-gzipped — this function is the deterministic
5556
alternative to ` + "`crane append --new_layer -`" + ` which re-gzipped

‎internal/nixstore/push.go‎

Lines changed: 80 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import (
1313

1414
"github.com/google/go-containerregistry/pkg/authn"
1515
"github.com/google/go-containerregistry/pkg/name"
16+
v1 "github.com/google/go-containerregistry/pkg/v1"
1617
"github.com/google/go-containerregistry/pkg/v1/mutate"
1718
"github.com/google/go-containerregistry/pkg/v1/remote"
1819
"github.com/google/go-containerregistry/pkg/v1/stream"
@@ -34,6 +35,11 @@ var ProgressTick = 5 * time.Second
3435
// disable chunking (single layer, original behavior).
3536
var ChunkSize int64 = 512 << 20 // 512 MB
3637

38+
// UploadJobs controls how many OCI layer blobs are uploaded concurrently
39+
// by remote.Write. Higher values improve throughput but risk GHCR
40+
// throttling; 16 caused session kills, 4 is a safe default.
41+
var UploadJobs = 4
42+
3743
// StallTimeout is the maximum duration Push tolerates zero byte
3844
// progress before aborting the upload. Set to 0 to disable. GHCR
3945
// occasionally throttles multi-GB layer uploads to near-zero; without
@@ -126,11 +132,11 @@ func fmtBytes(n uint64) string {
126132
// more OCI tar+gzip layers atop baseRef.
127133
//
128134
// When ChunkSize > 0 (default 512 MB), the input tar is split on entry
129-
// boundaries into ~ChunkSize chunks. Each chunk is pushed as a separate
130-
// OCI layer via its own remote.Write call — one upload at a time. After
131-
// each chunk the manifest is re-fetched so previous layers are referenced
132-
// by digest and skipped on subsequent pushes. No temp files; bytes flow
133-
// stdin → pipe → gzip → registry.
135+
// boundaries into ~ChunkSize temp files. All chunks are then uploaded
136+
// in parallel (UploadJobs concurrent blob streams, default 4) via a
137+
// single remote.Write call — one manifest write at the end. This
138+
// eliminates per-chunk manifest churn and overlaps gzip+upload across
139+
// layers. Temp files are removed after the push completes.
134140
//
135141
// When ChunkSize <= 0, the original single-layer streaming behavior
136142
// is preserved: bytes flow r → gzip → registry with no disk staging.
@@ -192,8 +198,9 @@ func Push(ctx context.Context, baseRef, dstRef string, r io.ReadCloser) error {
192198
if StallTimeout > 0 && cur > 0 && cur == lastBytes {
193199
ar := activeReader.Load()
194200
draining := ar != nil && ar.drained.Load()
195-
if draining {
196-
progressLog("[nix-store push] sent %s in %s — pipe drained, waiting for registry to finalize\n",
201+
allRead := TotalSizeHint > 0 && cur >= uint64(TotalSizeHint)
202+
if draining || allRead {
203+
progressLog("[nix-store push] sent %s in %s — all data read, waiting for registry to finalize\n",
197204
fmtBytes(cur), now.Sub(start).Round(time.Second))
198205
stallSince = time.Time{}
199206
} else if stallSince.IsZero() {
@@ -256,56 +263,91 @@ func Push(ctx context.Context, baseRef, dstRef string, r io.ReadCloser) error {
256263
if ChunkSize > 0 {
257264
defer r.Close()
258265
tr := tar.NewReader(r)
259-
currentImg := baseImg
260266

261-
uploadOpts := make([]remote.Option, len(remoteOpts))
262-
copy(uploadOpts, remoteOpts)
263-
uploadOpts = append(uploadOpts, remote.WithContext(uploadCtx))
267+
// Phase 1: split tar into temp files on disk.
268+
splitStart := time.Now()
269+
var chunkFiles []string
270+
defer func() {
271+
for _, f := range chunkFiles {
272+
os.Remove(f)
273+
}
274+
}()
264275

265276
for chunkIdx := 0; ; chunkIdx++ {
277+
tmpFile, err := os.CreateTemp("", fmt.Sprintf("nix-chunk-%03d-", chunkIdx))
278+
if err != nil {
279+
return fmt.Errorf("create temp chunk %d: %w", chunkIdx+1, err)
280+
}
281+
266282
pr, pw := io.Pipe()
267283
more := make(chan bool, 1)
268284
go writeOneChunk(tr, pw, ChunkSize, more)
269285

270-
progressLog("[nix-store push] chunk %d streaming\n", chunkIdx+1)
271-
272-
chunkPR := &progressReader{rc: pr, count: &byteCount}
273-
activeReader.Store(chunkPR)
286+
n, copyErr := io.Copy(tmpFile, pr)
287+
tmpFile.Close()
274288

275-
layer := stream.NewLayer(
276-
chunkPR,
277-
stream.WithMediaType(types.OCILayer),
278-
)
279-
280-
img, err := mutate.AppendLayers(currentImg, layer)
281-
if err != nil {
282-
pr.CloseWithError(err)
289+
if copyErr != nil {
290+
os.Remove(tmpFile.Name())
283291
<-more
284-
return fmt.Errorf("append chunk %d: %w", chunkIdx+1, err)
292+
return fmt.Errorf("buffer chunk %d: %w", chunkIdx+1, copyErr)
285293
}
286294

287-
writeStart := time.Now()
288-
if err := remote.Write(dst, img, uploadOpts...); err != nil {
289-
pr.CloseWithError(err)
290-
<-more
291-
if stallDetected.Load() {
292-
return fmt.Errorf("push %q: upload stalled — no progress for %s", dstRef, StallTimeout.Round(time.Second))
293-
}
294-
return fmt.Errorf("push chunk %d of %q: %w", chunkIdx+1, dstRef, err)
295-
}
296-
progressLog("[nix-store push] chunk %d committed (%s)\n", chunkIdx+1, time.Since(writeStart).Round(time.Millisecond))
295+
chunkFiles = append(chunkFiles, tmpFile.Name())
296+
progressLog("[nix-store push] chunk %d split (%s)\n", chunkIdx+1, fmtBytes(uint64(n)))
297297

298298
if !(<-more) {
299299
break
300300
}
301+
}
301302

302-
refetchStart := time.Now()
303-
currentImg, err = remote.Image(dst, uploadOpts...)
303+
var totalFileBytes int64
304+
for _, cf := range chunkFiles {
305+
if fi, err := os.Stat(cf); err == nil {
306+
totalFileBytes += fi.Size()
307+
}
308+
}
309+
TotalSizeHint = totalFileBytes
310+
progressLog("[nix-store push] %d chunks (%s) split in %s, uploading with %d parallel streams\n",
311+
len(chunkFiles), fmtBytes(uint64(totalFileBytes)), time.Since(splitStart).Round(time.Millisecond), UploadJobs)
312+
313+
// Phase 2: open all chunk files, wrap in layers, one remote.Write.
314+
var layers []v1.Layer
315+
var openFiles []*os.File
316+
defer func() {
317+
for _, f := range openFiles {
318+
f.Close()
319+
}
320+
}()
321+
322+
for _, chunkFile := range chunkFiles {
323+
f, err := os.Open(chunkFile)
304324
if err != nil {
305-
return fmt.Errorf("re-fetch after chunk %d: %w", chunkIdx+1, err)
325+
return fmt.Errorf("open chunk: %w", err)
326+
}
327+
openFiles = append(openFiles, f)
328+
329+
chunkPR := &progressReader{rc: f, count: &byteCount}
330+
layer := stream.NewLayer(chunkPR, stream.WithMediaType(types.OCILayer))
331+
layers = append(layers, layer)
332+
}
333+
334+
img, err := mutate.AppendLayers(baseImg, layers...)
335+
if err != nil {
336+
return fmt.Errorf("append %d layers: %w", len(layers), err)
337+
}
338+
339+
uploadOpts := make([]remote.Option, len(remoteOpts))
340+
copy(uploadOpts, remoteOpts)
341+
uploadOpts = append(uploadOpts, remote.WithContext(uploadCtx), remote.WithJobs(UploadJobs))
342+
343+
writeStart := time.Now()
344+
if err := remote.Write(dst, img, uploadOpts...); err != nil {
345+
if stallDetected.Load() {
346+
return fmt.Errorf("push %q: upload stalled — no progress for %s", dstRef, StallTimeout.Round(time.Second))
306347
}
307-
progressLog("[nix-store push] manifest re-fetched (%s)\n", time.Since(refetchStart).Round(time.Millisecond))
348+
return fmt.Errorf("push %q: %w", dstRef, err)
308349
}
350+
progressLog("[nix-store push] %d chunks committed (%s)\n", len(layers), time.Since(writeStart).Round(time.Millisecond))
309351
} else {
310352
pr := &progressReader{rc: r, count: &byteCount}
311353
stallCloser = pr

0 commit comments

Comments
 (0)