From 78993792a05c390c1bfbd4d854f55522986a162a Mon Sep 17 00:00:00 2001 From: Jan Kaluza Date: Thu, 13 Aug 2026 18:03:17 -0400 Subject: [PATCH 1/2] [podman-5.8] storage: wait for tar-split completion to avoid reader race Switch to `NewInputTarStreamWithDone` and explicitly wait for the tar-split goroutine to finish before returning. This prevents races with `uncompressed.Close()` while tar-split is still reading. Fixes: #148 Signed-off-by: Jan Kaluza (cherry picked from commit 91633032a252d7c54041f5590399eec781a16571) Signed-off-by: Qi Wang --- storage/go.mod | 2 +- storage/go.sum | 4 +- storage/layers.go | 21 +- storage/pkg/chunked/compression_linux_test.go | 4 +- storage/pkg/chunked/compressor/compressor.go | 276 +++++++++--------- storage/pkg/chunked/zstdchunked_test.go | 4 +- 6 files changed, 171 insertions(+), 140 deletions(-) diff --git a/storage/go.mod b/storage/go.mod index fec0003e84..bec999ac94 100644 --- a/storage/go.mod +++ b/storage/go.mod @@ -25,7 +25,7 @@ require ( github.com/stretchr/testify v1.11.1 github.com/tchap/go-patricia/v2 v2.3.3 github.com/ulikunitz/xz v0.5.15 - github.com/vbatts/tar-split v0.12.1 + github.com/vbatts/tar-split v0.12.3 golang.org/x/sync v0.17.0 golang.org/x/sys v0.37.0 gotest.tools/v3 v3.5.2 diff --git a/storage/go.sum b/storage/go.sum index 95b8c9d9de..edc7d1a859 100644 --- a/storage/go.sum +++ b/storage/go.sum @@ -71,8 +71,8 @@ github.com/tchap/go-patricia/v2 v2.3.3 h1:xfNEsODumaEcCcY3gI0hYPZ/PcpVv5ju6RMAhg github.com/tchap/go-patricia/v2 v2.3.3/go.mod h1:VZRHKAb53DLaG+nA9EaYYiaEx6YztwDlLElMsnSHD4k= github.com/ulikunitz/xz v0.5.15 h1:9DNdB5s+SgV3bQ2ApL10xRc35ck0DuIX/isZvIk+ubY= github.com/ulikunitz/xz v0.5.15/go.mod h1:nbz6k7qbPmH4IRqmfOplQw/tblSgqTqBwxkY0oWt/14= -github.com/vbatts/tar-split v0.12.1 h1:CqKoORW7BUWBe7UL/iqTVvkTBOF8UvOMKOIZykxnnbo= -github.com/vbatts/tar-split v0.12.1/go.mod h1:eF6B6i6ftWQcDqEn3/iGFRFRo8cBIMSJVOpnNdfTMFA= +github.com/vbatts/tar-split v0.12.3 h1:Cd46rkGXI3Td4yrVNwU8ripbxFaQbmesqhjBUUYAJSw= +github.com/vbatts/tar-split v0.12.3/go.mod h1:sQOc6OlqGCr7HkGx/IDBeKiTIvqhmj8KffNhEXG4Nq0= golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= diff --git a/storage/layers.go b/storage/layers.go index d485c9b4fa..33769dddac 100644 --- a/storage/layers.go +++ b/storage/layers.go @@ -2605,7 +2605,7 @@ func applyDiff(layerOptions *LayerOptions, diff io.Reader, tarSplitFile *os.File gidLog := make(map[uint32]struct{}) var uncompressedCounter *ioutils.WriteCounter - size, err := func() (int64, error) { // A scope for defer + size, err := func() (retSize int64, retErr error) { // A scope for defer compressor, err := pgzip.NewWriterLevel(tarSplitWriter, pgzip.BestSpeed) if err != nil { return -1, err @@ -2635,12 +2635,27 @@ func applyDiff(layerOptions *LayerOptions, diff io.Reader, tarSplitFile *os.File if uncompressedDigester != nil { uncompressedWriter = io.MultiWriter(uncompressedWriter, uncompressedDigester.Hash()) } - payload, err := asm.NewInputTarStream(io.TeeReader(uncompressed, uncompressedWriter), metadata, storage.NewDiscardFilePutter()) + payload, done, err := asm.NewInputTarStreamWithDone(io.TeeReader(uncompressed, uncompressedWriter), metadata, storage.NewDiscardFilePutter()) if err != nil { return -1, err } + defer func() { + payload.Close() + if doneErr := <-done; doneErr != nil && retErr == nil { + retErr = doneErr + } + }() - return applyDriverFunc(payload) + size, err := applyDriverFunc(payload) + if err != nil { + return -1, err + } + // Fully consume the payload; it may contain trailing zero padding, and we need all of that + // recorded in tar-split (which happens when the data passes through NewInputTarStreamWithDone). + if _, err := io.Copy(io.Discard, payload); err != nil { + return -1, err + } + return size, nil }() if err != nil { return nil, err diff --git a/storage/pkg/chunked/compression_linux_test.go b/storage/pkg/chunked/compression_linux_test.go index 5cae79dccd..183759f4bb 100644 --- a/storage/pkg/chunked/compression_linux_test.go +++ b/storage/pkg/chunked/compression_linux_test.go @@ -36,10 +36,12 @@ func TestTarSizeFromTarSplit(t *testing.T) { expectedTarSize := int64(tarball.Len()) var tarSplit bytes.Buffer - tsReader, err := asm.NewInputTarStream(&tarball, storage.NewJSONPacker(&tarSplit), storage.NewDiscardFilePutter()) + tsReader, done, err := asm.NewInputTarStreamWithDone(&tarball, storage.NewJSONPacker(&tarSplit), storage.NewDiscardFilePutter()) require.NoError(t, err) _, err = io.Copy(io.Discard, tsReader) require.NoError(t, err) + require.NoError(t, tsReader.Close()) + require.NoError(t, <-done) res, err := tarSizeFromTarSplit(&tarSplit) require.NoError(t, err) diff --git a/storage/pkg/chunked/compressor/compressor.go b/storage/pkg/chunked/compressor/compressor.go index ef26a812ba..68c4b834c6 100644 --- a/storage/pkg/chunked/compressor/compressor.go +++ b/storage/pkg/chunked/compressor/compressor.go @@ -240,173 +240,185 @@ func writeZstdChunkedStream(destFile io.Writer, outMetadata map[string]string, r } }() - its, err := asm.NewInputTarStream(reader, tarSplitData.packer, nil) - if err != nil { - return err - } + // Scope the NewInputTarStreamWithDone defer so we wait for done before + // returning to the outer function, which then closes tarSplitData.zstd. + metadata, err := func() (retMetadata []minimal.FileMetadata, retErr error) { + its, done, err := asm.NewInputTarStreamWithDone(reader, tarSplitData.packer, nil) + if err != nil { + return nil, err + } + defer func() { + its.Close() + if doneErr := <-done; doneErr != nil && retErr == nil { + retErr = doneErr + } + }() - tr := tar.NewReader(its) - tr.RawAccounting = true + tr := tar.NewReader(its) + tr.RawAccounting = true - buf := make([]byte, 4096) + buf := make([]byte, 4096) - zstdWriter, err := createZstdWriter(dest) - if err != nil { - return err - } - defer func() { - if zstdWriter != nil { - zstdWriter.Close() + zstdWriter, err := createZstdWriter(dest) + if err != nil { + return nil, err } - }() - - restartCompression := func() (int64, error) { - var offset int64 - if zstdWriter != nil { - if err := zstdWriter.Close(); err != nil { - return 0, err + defer func() { + if zstdWriter != nil { + zstdWriter.Close() } - offset = dest.Count - zstdWriter.Reset(dest) - } - return offset, nil - } + }() - var metadata []minimal.FileMetadata - for { - hdr, err := tr.Next() - if err != nil { - if err == io.EOF { - break + restartCompression := func() (int64, error) { + var offset int64 + if zstdWriter != nil { + if err := zstdWriter.Close(); err != nil { + return 0, err + } + offset = dest.Count + zstdWriter.Reset(dest) } - return err + return offset, nil } - rawBytes := tr.RawBytes() - if _, err := zstdWriter.Write(rawBytes); err != nil { - return err - } + var metadata []minimal.FileMetadata + for { + hdr, err := tr.Next() + if err != nil { + if err == io.EOF { + break + } + return nil, err + } - payloadDigester := digest.Canonical.Digester() - chunkDigester := digest.Canonical.Digester() + rawBytes := tr.RawBytes() + if _, err := zstdWriter.Write(rawBytes); err != nil { + return nil, err + } - // Now handle the payload, if any - startOffset := int64(0) - lastOffset := int64(0) - lastChunkOffset := int64(0) + payloadDigester := digest.Canonical.Digester() + chunkDigester := digest.Canonical.Digester() - checksum := "" + // Now handle the payload, if any + startOffset := int64(0) + lastOffset := int64(0) + lastChunkOffset := int64(0) - chunks := []chunk{} + checksum := "" - hf := &holesFinder{ - threshold: holesThreshold, - reader: bufio.NewReader(tr), - } + chunks := []chunk{} - rcReader := &rollingChecksumReader{ - reader: hf, - rollsum: NewRollSum(), - } + hf := &holesFinder{ + threshold: holesThreshold, + reader: bufio.NewReader(tr), + } - payloadDest := io.MultiWriter(payloadDigester.Hash(), chunkDigester.Hash(), zstdWriter) - for { - mustSplit, read, errRead := rcReader.Read(buf) - if errRead != nil && errRead != io.EOF { - return err + rcReader := &rollingChecksumReader{ + reader: hf, + rollsum: NewRollSum(), } - // restart the compression only if there is a payload. - if read > 0 { - if startOffset == 0 { - startOffset, err = restartCompression() - if err != nil { - return err - } - lastOffset = startOffset - } - if _, err := payloadDest.Write(buf[:read]); err != nil { - return err + payloadDest := io.MultiWriter(payloadDigester.Hash(), chunkDigester.Hash(), zstdWriter) + for { + mustSplit, read, errRead := rcReader.Read(buf) + if errRead != nil && errRead != io.EOF { + return nil, errRead } - } - if (mustSplit || errRead == io.EOF) && startOffset > 0 { - off, err := restartCompression() - if err != nil { - return err + // restart the compression only if there is a payload. + if read > 0 { + if startOffset == 0 { + startOffset, err = restartCompression() + if err != nil { + return nil, err + } + lastOffset = startOffset + } + + if _, err := payloadDest.Write(buf[:read]); err != nil { + return nil, err + } } + if (mustSplit || errRead == io.EOF) && startOffset > 0 { + off, err := restartCompression() + if err != nil { + return nil, err + } - chunkSize := rcReader.WrittenOut - lastChunkOffset - if chunkSize > 0 { - chunkType := minimal.ChunkTypeData - if rcReader.IsLastChunkZeros { - chunkType = minimal.ChunkTypeZeros + chunkSize := rcReader.WrittenOut - lastChunkOffset + if chunkSize > 0 { + chunkType := minimal.ChunkTypeData + if rcReader.IsLastChunkZeros { + chunkType = minimal.ChunkTypeZeros + } + + chunks = append(chunks, chunk{ + ChunkOffset: lastChunkOffset, + Offset: lastOffset, + Checksum: chunkDigester.Digest().String(), + ChunkSize: chunkSize, + ChunkType: chunkType, + }) } - chunks = append(chunks, chunk{ - ChunkOffset: lastChunkOffset, - Offset: lastOffset, - Checksum: chunkDigester.Digest().String(), - ChunkSize: chunkSize, - ChunkType: chunkType, - }) + lastOffset = off + lastChunkOffset = rcReader.WrittenOut + chunkDigester = digest.Canonical.Digester() + payloadDest = io.MultiWriter(payloadDigester.Hash(), chunkDigester.Hash(), zstdWriter) + } + if errRead == io.EOF { + if startOffset > 0 { + checksum = payloadDigester.Digest().String() + } + break } + } - lastOffset = off - lastChunkOffset = rcReader.WrittenOut - chunkDigester = digest.Canonical.Digester() - payloadDest = io.MultiWriter(payloadDigester.Hash(), chunkDigester.Hash(), zstdWriter) + mainEntry, err := minimal.NewFileMetadata(hdr) + if err != nil { + return nil, err } - if errRead == io.EOF { - if startOffset > 0 { - checksum = payloadDigester.Digest().String() + mainEntry.Digest = checksum + mainEntry.Offset = startOffset + mainEntry.EndOffset = lastOffset + entries := []minimal.FileMetadata{mainEntry} + for i := 1; i < len(chunks); i++ { + entries = append(entries, minimal.FileMetadata{ + Type: minimal.TypeChunk, + Name: hdr.Name, + ChunkOffset: chunks[i].ChunkOffset, + }) + } + if len(chunks) > 1 { + for i := range chunks { + entries[i].ChunkSize = chunks[i].ChunkSize + entries[i].Offset = chunks[i].Offset + entries[i].ChunkDigest = chunks[i].Checksum + entries[i].ChunkType = chunks[i].ChunkType } - break } + metadata = append(metadata, entries...) } - mainEntry, err := minimal.NewFileMetadata(hdr) - if err != nil { - return err - } - mainEntry.Digest = checksum - mainEntry.Offset = startOffset - mainEntry.EndOffset = lastOffset - entries := []minimal.FileMetadata{mainEntry} - for i := 1; i < len(chunks); i++ { - entries = append(entries, minimal.FileMetadata{ - Type: minimal.TypeChunk, - Name: hdr.Name, - ChunkOffset: chunks[i].ChunkOffset, - }) - } - if len(chunks) > 1 { - for i := range chunks { - entries[i].ChunkSize = chunks[i].ChunkSize - entries[i].Offset = chunks[i].Offset - entries[i].ChunkDigest = chunks[i].Checksum - entries[i].ChunkType = chunks[i].ChunkType - } + rawBytes := tr.RawBytes() + if _, err := zstdWriter.Write(rawBytes); err != nil { + return nil, err } - metadata = append(metadata, entries...) - } - rawBytes := tr.RawBytes() - if _, err := zstdWriter.Write(rawBytes); err != nil { - zstdWriter.Close() - return err - } - - // make sure the entire tarball is flushed to the output as it might contain - // some trailing zeros that affect the checksum. - if _, err := io.Copy(zstdWriter, its); err != nil { - zstdWriter.Close() - return err - } + // make sure the entire tarball is flushed to the output as it might contain + // some trailing zeros that affect the checksum. + if _, err := io.Copy(zstdWriter, its); err != nil { + return nil, err + } - if err := zstdWriter.Close(); err != nil { + if err := zstdWriter.Close(); err != nil { + return nil, err + } + zstdWriter = nil + return metadata, nil + }() + if err != nil { return err } - zstdWriter = nil if err := tarSplitData.zstd.Close(); err != nil { return err diff --git a/storage/pkg/chunked/zstdchunked_test.go b/storage/pkg/chunked/zstdchunked_test.go index 435342c2c1..2a65ffcf5d 100644 --- a/storage/pkg/chunked/zstdchunked_test.go +++ b/storage/pkg/chunked/zstdchunked_test.go @@ -109,10 +109,12 @@ func TestGenerateAndParseManifest(t *testing.T) { err := tsTarW.Close() require.NoError(t, err) var tarSplitUncompressed bytes.Buffer - tsReader, err := asm.NewInputTarStream(&tsTarball, storage.NewJSONPacker(&tarSplitUncompressed), storage.NewDiscardFilePutter()) + tsReader, done, err := asm.NewInputTarStreamWithDone(&tsTarball, storage.NewJSONPacker(&tarSplitUncompressed), storage.NewDiscardFilePutter()) require.NoError(t, err) _, err = io.Copy(io.Discard, tsReader) require.NoError(t, err) + require.NoError(t, tsReader.Close()) + require.NoError(t, <-done) encoder, err := zstd.NewWriter(nil) if err != nil { From cc3e1ba032f971fb27749766fcc0d1461424abc6 Mon Sep 17 00:00:00 2001 From: Paul Holzinger Date: Mon, 27 Jul 2026 11:34:18 +0200 Subject: [PATCH 2/2] ci: use larger runners For some yet unknown reasons the performance of the runners is a lot worse than before causing timeouts for most tasks. Use some bigger runner to try to work around. I added a comment on the cncf ci infra channel, it seems we are not the only one with problems. Signed-off-by: Paul Holzinger Signed-off-by: Qi Wang --- .github/workflows/ci.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 26d40c2e91..80479462a6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -52,7 +52,7 @@ jobs: driver: overlay-transient uses: ./.github/workflows/lima.yml with: - runner: cncf-ubuntu-2-8-x86 + runner: cncf-ubuntu-4-16-x86 module: storage distro: ${{ matrix.distro }} variant: ${{ matrix.driver }} @@ -114,7 +114,7 @@ jobs: module: [image, image-skopeo] uses: ./.github/workflows/lima.yml with: - runner: cncf-ubuntu-2-8-x86 + runner: cncf-ubuntu-4-16-x86 module: ${{ matrix.module }} distro: fedora-current variant: ${{ matrix.variant }} @@ -129,7 +129,7 @@ jobs: name: "common fedora-current" uses: ./.github/workflows/lima.yml with: - runner: cncf-ubuntu-2-8-x86 + runner: cncf-ubuntu-4-16-x86 module: common distro: fedora-current timeout: 20