diff --git a/docs/adopting-packstore.md b/docs/adopting-packstore.md index ac7d7a87..c503be64 100644 --- a/docs/adopting-packstore.md +++ b/docs/adopting-packstore.md @@ -193,7 +193,8 @@ required. Opt into bounded-memory writing with `Writer.AppendStream`, or prepare entries concurrently with `PrepareBlob` and append the resulting `PreparedBlob` values in deterministic order with `AppendPrepared`. -Use `Reader.OpenBlob` for bounded-memory plain reads. The supplied `Entry` must +Use `Reader.OpenBlob` for streaming plain reads. Legacy single-segment frames +can still require a decoder window as large as the blob. The supplied `Entry` must come from that reader's immutable `Entries` result; caller-constructed offsets are rejected. `BlobReader` has the same terminal verification and early-close contract as `Store.OpenStream`. Use `ReaderOptions` and `BlobReaderOptions` to @@ -262,9 +263,13 @@ of exact-owned staging files in recovery. `ReaderSlots` bounds idle cached descriptors, not all live descriptors. Each active stream holds a lease even after eviction, so budget file descriptors as `ReaderSlots + maximum concurrent streams`, plus unrelated application use. -Compressed reads also enforce a decoder-window limit. A legacy format-v1 frame -whose declared window exceeds policy fails with a typed limit error rather -than increasing memory implicitly. +Packed `OpenStream` and `ReadBounded` reads also enforce a decoder-window limit +derived from the blob policy. A legacy format-v1 frame whose declared window +exceeds that limit fails with a typed limit error. Direct `pack` readers with +zero `ReaderLimits.WindowBytes` allow windows up to 4 GiB on 64-bit systems, +so a legacy single-segment frame can require memory proportional to its full +raw size. The decoder window ceiling remains 512 MiB on 32-bit systems. +The buffered compatibility path `Store.Open` uses that default as well. Format v1 represents raw lengths up to 4 GiB. Its encrypted frames use whole-entry authenticated encryption and cannot safely expose streamed diff --git a/docs/architecture/streaming-pack-io.md b/docs/architecture/streaming-pack-io.md index b38e4d06..ea0c82dd 100644 --- a/docs/architecture/streaming-pack-io.md +++ b/docs/architecture/streaming-pack-io.md @@ -1,9 +1,9 @@ # Streaming pack I/O -Kit reads and writes format-v1 packs without requiring memory proportional to -the largest object. This document describes the contracts and ownership rules -that keep the implementation bounded, verifiable, and safe during concurrent -maintenance. +Kit streams format-v1 packs with memory bounded by buffers, decoder windows, +and concurrency. Legacy single-segment frames can require a decoder window as +large as the object. This document describes the contracts and ownership rules +for verification and concurrent maintenance. The on-disk representation remains the format described by the [backup repository and pack format](../../backup/FORMAT.md). Streaming changes @@ -13,8 +13,7 @@ how bytes move through the implementation, not how existing packs are encoded. Streaming pack I/O is designed to: -- bound heap use by configured buffers, decoder windows, and concurrency rather - than object size; +- bound heap use by configured buffers, decoder windows, and concurrency; - preserve format-v1 compatibility and the existing buffered APIs; - verify stored and decoded bytes before treating output as authoritative; - keep resource-policy limits explicit and independently configurable; and @@ -158,9 +157,15 @@ Resource limits describe distinct risks and should remain distinct: - scratch-space limits bound preparation and staging disk use; and - concurrency limits bound aggregate buffers and open descriptors. -Compressed entries use a bounded decoder window. Legacy format-v1 entries may -declare a window related to their raw size; when that exceeds policy, the read -fails with a typed policy error instead of silently increasing memory use. +With default `pack.ReaderLimits`, compressed streams allow decoder windows up to +`pack.MaxRawLen` (4 GiB) on 64-bit systems. On 32-bit systems, both buffered and +streaming decoders retain a 512 MiB window ceiling to avoid codec integer +overflow. Legacy single-segment frames can require a window as large as the +blob, so their streaming reads can allocate that much memory. +Set `ReaderLimits.WindowBytes` or a narrower `BlobReaderOptions.WindowBytes` +to impose a lower ceiling; a stream exceeding it fails with a typed policy error. +These window limits do not apply to buffered `Reader.ReadBlob` calls. +New buffered writes over 512 MiB use an explicit, smaller window. Format-v1 encrypted entries remain buffered because their whole-entry authenticated-encryption construction cannot authenticate a streamed prefix. @@ -168,8 +173,9 @@ True encrypted streaming requires a chunk-authenticated representation in a future format. The format-v1 length fields also impose their existing object-size ceiling. -Streaming removes heap proportionality within that ceiling; it does not extend -the wire format beyond it. +Streaming avoids whole-object output buffers, but legacy decoder windows can +still require memory proportional to object size. It does not extend the wire +format's size ceiling. ## Implementation evidence diff --git a/pack/frame.go b/pack/frame.go index 15e80c19..173ce16e 100644 --- a/pack/frame.go +++ b/pack/frame.go @@ -2,19 +2,30 @@ package pack import ( "fmt" + "strconv" "sync" "github.com/klauspost/compress/zstd" ) +type zstdEncoderKey struct { + level int + singleSegment bool +} + // Shared zstd codecs. EncodeAll/DecodeAll are safe for concurrent use on a -// single Encoder/Decoder, so one instance per level serves all writers. +// single Encoder/Decoder, so one instance per level and framing mode serves +// all writers. var ( zstdEncMu sync.Mutex - zstdEncs = map[int]*zstd.Encoder{} + zstdEncs = map[zstdEncoderKey]*zstd.Encoder{} zstdDec = func() *zstd.Decoder { d, err := zstd.NewReader(nil, - zstd.WithDecoderConcurrency(0), zstd.WithDecoderMaxMemory(1<<32)) + zstd.WithDecoderConcurrency(0), zstd.WithDecoderMaxMemory(1<<32), + // Single-segment frames carry a window equal to their content + // size; the 512 MiB default would reject blobs this package + // itself produced. + zstd.WithDecoderMaxWindow(maxDecoderWindow())) if err != nil { panic(fmt.Sprintf("pack: initializing zstd decoder: %v", err)) } @@ -22,22 +33,32 @@ var ( }() ) -func zstdEncoder(level int) *zstd.Encoder { +func maxDecoderWindow() uint64 { + // The codec uses int for history-buffer sizes and adds space beyond the + // window. Retain its default on 32-bit systems to avoid integer overflow. + if strconv.IntSize == 32 { + return zstd.MaxWindowSize + } + return MaxRawLen +} + +func zstdEncoder(level int, singleSegment bool) *zstd.Encoder { if level <= 0 { level = DefaultZstdLevel } + key := zstdEncoderKey{level: level, singleSegment: singleSegment} zstdEncMu.Lock() defer zstdEncMu.Unlock() - if enc, ok := zstdEncs[level]; ok { + if enc, ok := zstdEncs[key]; ok { return enc } enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.EncoderLevelFromZstd(level)), - zstd.WithSingleSegment(true)) + zstd.WithSingleSegment(singleSegment)) if err != nil { panic(fmt.Sprintf("pack: initializing zstd encoder level %d: %v", level, err)) } - zstdEncs[level] = enc + zstdEncs[key] = enc return enc } @@ -55,12 +76,15 @@ func minCompressionSavings(rawLen int) int { func encodeFrame(raw []byte, level int) (stored []byte, compressed bool) { // Legacy bounded readers cap decoder memory at RawLen. A zstd frame always // needs at least MinWindowSize, so smaller blobs must remain raw for - // downgrade compatibility. Single-segment frames keep larger windows tied - // to the authoritative content length. + // downgrade compatibility. if len(raw) < zstd.MinWindowSize { return raw, false } - c := zstdEncoder(level).EncodeAll(raw, make([]byte, 0, len(raw))) + // Single-segment frames use their content size as the window. Above the + // klauspost decoder's default 512 MiB limit, use an explicit window instead. + singleSegment := uint64(len(raw)) <= uint64(zstd.MaxWindowSize) + encoder := zstdEncoder(level, singleSegment) + c := encoder.EncodeAll(raw, make([]byte, 0, len(raw))) minSavings := minCompressionSavings(len(raw)) if len(c) > len(raw)-minSavings { return raw, false diff --git a/pack/oversized_frame_test.go b/pack/oversized_frame_test.go new file mode 100644 index 00000000..96900fbf --- /dev/null +++ b/pack/oversized_frame_test.go @@ -0,0 +1,108 @@ +//go:build !race + +package pack + +import ( + "crypto/sha256" + "io" + "path/filepath" + "strconv" + "testing" + + "github.com/klauspost/compress/zstd" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// These tests need real blobs over 512 MiB and a 64-bit address space. Run them +// in ordinary CI, but omit race builds where instrumentation multiplies their +// memory cost. +func TestEncodeFrameRoundTripsBlobsOverMaxWindow(t *testing.T) { + if testing.Short() || strconv.IntSize == 32 { + t.Skip("requires a 64-bit process and more than 512 MiB") + } + + const size = zstd.MaxWindowSize + 1 + raw := make([]byte, size) + for i := range raw { + raw[i] = byte(i % 251) + } + wantHash := sha256.Sum256(raw) + stored, compressed := EncodeFrame(raw, DefaultZstdLevel) + require.True(t, compressed) + + var header zstd.Header + require.NoError(t, header.Decode(stored)) + require.False(t, header.SingleSegment, + "blobs over the window limit must not use single-segment framing") + assert.LessOrEqual(t, header.WindowSize, uint64(zstd.MaxWindowSize)) + + // Keep the previous reader's 512 MiB window limit. Reuse the input buffer + // after encoding to avoid another allocation proportional to the blob. + decoder, err := zstd.NewReader(nil, + zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(MaxRawLen)) + require.NoError(t, err) + t.Cleanup(decoder.Close) + decoded, err := decoder.DecodeAll(stored, raw[:0]) + require.NoError(t, err) + assert.Len(t, decoded, size) + assert.Equal(t, wantHash, sha256.Sum256(decoded)) +} + +func TestReaderReadsOversizedSingleSegmentFrame(t *testing.T) { + if testing.Short() || strconv.IntSize == 32 { + t.Skip("requires a 64-bit process and more than 512 MiB") + } + + const size = zstd.MaxWindowSize + 1 + raw := make([]byte, size) + for i := range raw { + raw[i] = byte(i % 251) + } + wantHash := sha256.Sum256(raw) + + // Reproduce the old writer independently of EncodeFrame, which now avoids + // single-segment framing for oversized blobs. + encoder, err := zstd.NewWriter(nil, + zstd.WithEncoderLevel(zstd.EncoderLevelFromZstd(DefaultZstdLevel)), + zstd.WithEncoderConcurrency(1), zstd.WithSingleSegment(true)) + require.NoError(t, err) + stored := encoder.EncodeAll(raw, nil) + require.NoError(t, encoder.Close()) + var header zstd.Header + require.NoError(t, header.Decode(stored)) + require.True(t, header.SingleSegment) + require.Equal(t, uint64(size), header.FrameContentSize) + + dir := t.TempDir() + writer, err := NewWriter(dir, WriterOptions{}) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, writer.Abort()) }) + _, err = writer.AppendEncoded(BlobID(wantHash), stored, size, true) + require.NoError(t, err) + path := filepath.Join(dir, writer.ID()+".pack") + _, err = writer.Seal(path) + require.NoError(t, err) + reader, err := OpenReader(path, nil) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, reader.Close()) }) + entry := reader.Entries()[0] + + t.Run("ReadBlob", func(t *testing.T) { + decoded, err := reader.ReadBlob(entry) + require.NoError(t, err) + assert.Len(t, decoded, size) + assert.Equal(t, wantHash, sha256.Sum256(decoded)) + }) + t.Run("OpenBlob", func(t *testing.T) { + stream, err := reader.OpenBlob(t.Context(), entry) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, stream.Close()) }) + hash := sha256.New() + n, err := io.Copy(hash, stream) + require.NoError(t, err) + assert.Equal(t, int64(size), n) + assert.Equal(t, wantHash[:], hash.Sum(nil)) + assert.True(t, stream.Verified()) + }) +} diff --git a/pack/reader.go b/pack/reader.go index 6f978983..4a9ce459 100644 --- a/pack/reader.go +++ b/pack/reader.go @@ -11,19 +11,23 @@ import ( "path/filepath" "strings" "sync" - - "github.com/klauspost/compress/zstd" ) // ReaderLimits bounds format-controlled quantities before allocation or -// exposure. Zero fields select the corresponding format maximum. +// exposure. Zero fields select the corresponding format maximum, except for +// the decoder window's lower ceiling on 32-bit systems. type ReaderLimits struct { ContainerBytes uint64 FooterBytes uint64 Entries uint64 RawBytes uint64 StoredBytes uint64 - WindowBytes uint64 + + // WindowBytes limits each streaming decoder's window. Zero allows up to + // MaxRawLen (4 GiB) on 64-bit systems, or 512 MiB on 32-bit systems. Old + // single-segment frames can require their full raw size as a window. + // Set a smaller limit to bound that memory use. + WindowBytes uint64 } // ReaderOptions configures bounded pack opening. @@ -174,7 +178,13 @@ func normalizeReaderLimits(limits ReaderLimits) ReaderLimits { limits.Entries = set(limits.Entries, MaxFooterLen/entrySize) limits.RawBytes = set(limits.RawBytes, MaxRawLen) limits.StoredBytes = set(limits.StoredBytes, MaxStoredLen) - limits.WindowBytes = set(limits.WindowBytes, zstd.MaxWindowSize) + // A single-segment frame declares no window, so its effective window is the + // whole frame content size. Capping the default at zstd.MaxWindowSize + // (512 MiB) rejected larger single-segment blobs this package wrote itself. + // The format ceiling allows those reads at the cost of a potentially + // blob-sized decoder window. OpenBlobWithOptions enforces explicit lower + // window limits before creating the decoder. + limits.WindowBytes = set(limits.WindowBytes, maxDecoderWindow()) return limits } diff --git a/pack/reader_test.go b/pack/reader_test.go index 220a0e7e..44c7d38c 100644 --- a/pack/reader_test.go +++ b/pack/reader_test.go @@ -3,15 +3,59 @@ package pack import ( "bytes" "crypto/rand" + "math" "os" "path/filepath" + "strconv" "strings" "testing" + "github.com/klauspost/compress/zstd" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) +func TestReaderRejectsOversized32BitWindow(t *testing.T) { + if strconv.IntSize != 32 { + t.Skip("32-bit decoder window limit") + } + // Only the header and an empty final block are needed: the window must be + // rejected before the codec calculates an overflowing history-buffer size. + header := zstd.Header{SingleSegment: true, HasFCS: true, FrameContentSize: math.MaxInt32} + frame, err := header.AppendTo(nil) + require.NoError(t, err) + frame = append(frame, 1, 0, 0) + dir := t.TempDir() + writer, err := NewWriter(dir, WriterOptions{}) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, writer.Abort()) }) + _, err = writer.AppendEncoded(BlobID{}, frame, math.MaxInt32, true) + require.NoError(t, err) + path := filepath.Join(dir, writer.ID()+".pack") + _, err = writer.Seal(path) + require.NoError(t, err) + reader, err := OpenReader(path, nil) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, reader.Close()) }) + entry := reader.Entries()[0] + + t.Run("ReadBlob", func(t *testing.T) { + _, err := reader.ReadBlob(entry) + require.ErrorIs(t, err, zstd.ErrWindowSizeExceeded) + }) + t.Run("OpenBlob", func(t *testing.T) { + stream, err := reader.OpenBlob(t.Context(), entry) + if stream != nil { + t.Cleanup(func() { _ = stream.Close() }) + } + var limitErr *StreamLimitError + require.ErrorAs(t, err, &limitErr) + assert.Equal(t, StreamLimitWindowBytes, limitErr.Dimension) + assert.Equal(t, uint64(math.MaxInt32), limitErr.Actual) + assert.Equal(t, uint64(512<<20), limitErr.Limit) + }) +} + // buildTestPack writes a pack with the given blobs and returns its final path // and entries. crypter may be nil for a plain pack. func buildTestPack(t *testing.T, blobs [][]byte,