diff --git a/common/go.mod b/common/go.mod index ff62daf579..297068fafb 100644 --- a/common/go.mod +++ b/common/go.mod @@ -61,6 +61,7 @@ require ( github.com/VividCortex/ewma v1.2.0 // indirect github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d // indirect github.com/chzyer/readline v1.5.1 // indirect + github.com/clipperhouse/uax29/v2 v2.7.0 // indirect github.com/containerd/errdefs v1.0.0 // indirect github.com/containerd/errdefs/pkg v0.3.0 // indirect github.com/containerd/log v0.1.0 // indirect @@ -90,7 +91,7 @@ require ( github.com/kr/fs v0.1.0 // indirect github.com/letsencrypt/boulder v0.0.0-20240620165639-de9c06129bec // indirect github.com/manifoldco/promptui v0.9.0 // indirect - github.com/mattn/go-runewidth v0.0.16 // indirect + github.com/mattn/go-runewidth v0.0.20 // indirect github.com/mattn/go-sqlite3 v1.14.32 // indirect github.com/miekg/pkcs11 v1.1.1 // indirect github.com/mistifyio/go-zfs/v3 v3.1.0 // indirect @@ -102,7 +103,6 @@ require ( github.com/modern-go/reflect2 v1.0.2 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/proglottis/gpgme v0.1.5 // indirect - github.com/rivo/uniseg v0.4.7 // indirect github.com/secure-systems-lab/go-securesystemslib v0.9.1 // indirect github.com/sigstore/fulcio v1.7.1 // indirect github.com/sigstore/protobuf-specs v0.4.1 // indirect @@ -113,7 +113,7 @@ require ( github.com/tchap/go-patricia/v2 v2.3.3 // indirect github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399 // indirect github.com/ulikunitz/xz v0.5.15 // indirect - github.com/vbatts/tar-split v0.12.1 // indirect + github.com/vbatts/tar-split v0.12.3 // indirect github.com/vbauerster/mpb/v8 v8.10.2 // indirect github.com/vishvananda/netns v0.0.5 // indirect github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb // indirect diff --git a/common/go.sum b/common/go.sum index e304ba542f..f9ebb0bc76 100644 --- a/common/go.sum +++ b/common/go.sum @@ -33,6 +33,8 @@ github.com/chzyer/readline v1.5.1/go.mod h1:Eh+b79XXUwfKfcPLepksvw2tcLE/Ct21YObk github.com/chzyer/test v0.0.0-20180213035817-a1ea475d72b1/go.mod h1:Q3SI9o4m/ZMnBNeIyt5eFwwo7qiLfzFZmjNmxjkiQlU= github.com/chzyer/test v1.0.0 h1:p3BQDXSxOhOG0P9z6/hGnII4LGiEPOYBhs8asl/fC04= github.com/chzyer/test v1.0.0/go.mod h1:2JlltgoNkt4TW/z9V/IzDdFaMTM2JPIi26O1pF38GC8= +github.com/clipperhouse/uax29/v2 v2.7.0 h1:+gs4oBZ2gPfVrKPthwbMzWZDaAFPGYK72F0NJv2v7Vk= +github.com/clipperhouse/uax29/v2 v2.7.0/go.mod h1:EFJ2TJMRUaplDxHKj1qAEhCtQPW2tJSwu5BF98AuoVM= github.com/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= @@ -155,8 +157,8 @@ github.com/manifoldco/promptui v0.9.0 h1:3V4HzJk1TtXW1MTZMP7mdlwbBpIinw3HztaIlYt github.com/manifoldco/promptui v0.9.0/go.mod h1:ka04sppxSGFAtxX0qhlYQjISsg9mR4GWtQEhdbn6Pgg= github.com/maruel/natural v1.1.1 h1:Hja7XhhmvEFhcByqDoHz9QZbkWey+COd9xWfCfn1ioo= github.com/maruel/natural v1.1.1/go.mod h1:v+Rfd79xlw1AgVBjbO0BEQmptqb5HvL/k9GRHB7ZKEg= -github.com/mattn/go-runewidth v0.0.16 h1:E5ScNMtiwvlvB5paMFdw9p4kSQzbXFikJ5SQO6TULQc= -github.com/mattn/go-runewidth v0.0.16/go.mod h1:Jdepj2loyihRzMpdS35Xk/zdY8IAYHsh153qUoGf23w= +github.com/mattn/go-runewidth v0.0.20 h1:WcT52H91ZUAwy8+HUkdM3THM6gXqXuLJi9O3rjcQQaQ= +github.com/mattn/go-runewidth v0.0.20/go.mod h1:XBkDxAl56ILZc9knddidhrOlY5R/pDhgLpndooCuJAs= github.com/mattn/go-sqlite3 v1.14.32 h1:JD12Ag3oLy1zQA+BNn74xRgaBbdhbNIDYvQUEuuErjs= github.com/mattn/go-sqlite3 v1.14.32/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= github.com/mfridman/tparse v0.18.0 h1:wh6dzOKaIwkUGyKgOntDW4liXSo37qg5AXbIhkMV3vE= @@ -227,9 +229,6 @@ github.com/prometheus/common v0.63.0 h1:YR/EIY1o3mEFP/kZCD7iDMnLPlGyuU2Gb3HIcXnA github.com/prometheus/common v0.63.0/go.mod h1:VVFF/fBIoToEnWRVkYoXEkq3R3paCoxG9PXP74SnV18= github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= -github.com/rivo/uniseg v0.2.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc= -github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ= -github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII= github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= @@ -286,8 +285,8 @@ github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399 h1:e/5i7d4oYZ+C github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399/go.mod h1:LdwHTNJT99C5fTAzDz0ud328OgXz+gierycbcIx2fRs= 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= github.com/vbauerster/mpb/v8 v8.10.2 h1:2uBykSHAYHekE11YvJhKxYmLATKHAGorZwFlyNw4hHM= github.com/vbauerster/mpb/v8 v8.10.2/go.mod h1:+Ja4P92E3/CorSZgfDtK46D7AVbDqmBQRTmyTqPElo0= github.com/vishvananda/netlink v1.3.1 h1:3AEMt62VKqz90r0tmNhog0r/PpWKmrEShJU0wJW6bV0= diff --git a/image/go.mod b/image/go.mod index bc0d2361c8..1266ee0996 100644 --- a/image/go.mod +++ b/image/go.mod @@ -1,6 +1,6 @@ module go.podman.io/image/v5 -go 1.24.0 +go 1.24.2 // Warning: Ensure the "go" and "toolchain" versions match exactly to prevent unwanted auto-updates. // That generally means there should be no toolchain directive present. @@ -62,7 +62,7 @@ require ( github.com/docker/go-units v0.5.0 // indirect github.com/docker/libtrust v0.0.0-20160708172513-aabc10ec26b7 // indirect github.com/felixge/httpsnoop v1.0.4 // indirect - github.com/go-jose/go-jose/v4 v4.0.5 // indirect + github.com/go-jose/go-jose/v4 v4.1.3 // indirect github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/golang/protobuf v1.5.4 // indirect @@ -101,7 +101,7 @@ require ( github.com/stefanberger/go-pkcs11uri v0.0.0-20230803200340-78284954bff6 // indirect github.com/tchap/go-patricia/v2 v2.3.3 // indirect github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399 // indirect - github.com/vbatts/tar-split v0.12.1 // indirect + github.com/vbatts/tar-split v0.12.3 // indirect go.opentelemetry.io/auto/sdk v1.1.0 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 // indirect go.opentelemetry.io/otel v1.36.0 // indirect diff --git a/image/go.sum b/image/go.sum index f107d021dd..dc5bd402da 100644 --- a/image/go.sum +++ b/image/go.sum @@ -75,8 +75,8 @@ github.com/fatih/color v1.16.0 h1:zmkK9Ngbjj+K0yRhTVONQh1p/HknKYSlNT+vZCzyokM= github.com/fatih/color v1.16.0/go.mod h1:fL2Sau1YI5c0pdGEVCbKQbLXB6edEj1ZgiY4NijnWvE= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= -github.com/go-jose/go-jose/v4 v4.0.5 h1:M6T8+mKZl/+fNNuFHvGIzDz7BTLQPIounk/b9dw3AaE= -github.com/go-jose/go-jose/v4 v4.0.5/go.mod h1:s3P1lRrkT8igV8D9OjyL4WRyHvjB6a4JSllnOrmmBOA= +github.com/go-jose/go-jose/v4 v4.1.3 h1:CVLmWDhDVRa6Mi/IgCgaopNosCaHz7zrMeF9MlZRkrs= +github.com/go-jose/go-jose/v4 v4.1.3/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE= github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V4qmtdjCk= @@ -263,8 +263,8 @@ github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399 h1:e/5i7d4oYZ+C github.com/titanous/rocacheck v0.0.0-20171023193734-afe73141d399/go.mod h1:LdwHTNJT99C5fTAzDz0ud328OgXz+gierycbcIx2fRs= 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= github.com/vbauerster/mpb/v8 v8.10.2 h1:2uBykSHAYHekE11YvJhKxYmLATKHAGorZwFlyNw4hHM= github.com/vbauerster/mpb/v8 v8.10.2/go.mod h1:+Ja4P92E3/CorSZgfDtK46D7AVbDqmBQRTmyTqPElo0= github.com/ysmood/fetchup v0.2.3 h1:ulX+SonA0Vma5zUFXtv52Kzip/xe7aj4vqT5AJwQ+ZQ= diff --git a/storage/go.mod b/storage/go.mod index fec0003e84..9b1cef5845 100644 --- a/storage/go.mod +++ b/storage/go.mod @@ -1,4 +1,4 @@ -go 1.24.0 +go 1.24.2 // Warning: Ensure the "go" and "toolchain" versions match exactly to prevent unwanted auto-updates. // That generally means there should be no toolchain directive present. @@ -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 {