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
13 changes: 9 additions & 4 deletions docs/adopting-packstore.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
28 changes: 17 additions & 11 deletions docs/architecture/streaming-pack-io.md
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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
Expand Down Expand Up @@ -158,18 +157,25 @@ 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.
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

Expand Down
44 changes: 34 additions & 10 deletions pack/frame.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,42 +2,63 @@ 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))
}
return d
}()
)

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
}

Expand All @@ -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
Expand Down
108 changes: 108 additions & 0 deletions pack/oversized_frame_test.go
Original file line number Diff line number Diff line change
@@ -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())
})
}
20 changes: 15 additions & 5 deletions pack/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
}

Expand Down
44 changes: 44 additions & 0 deletions pack/reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down