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
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ func (s *fakeSeekable) Size(_ context.Context) (int64, error) {
return int64(len(s.data)), nil
}

func (s *fakeSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FrameTable, [32]byte, error) {
func (s *fakeSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FullFrameTable, [32]byte, error) {
panic("not used")
}

Expand Down Expand Up @@ -146,7 +146,7 @@ func makeCompressedTestData(tb testing.TB, data []byte) (*storage.FrameTable, *f
})
require.NoError(tb, err)

return ft, &fakeSeekable{data: compressed}
return ft.Table(), &fakeSeekable{data: compressed}
}

type chunkerTestCase struct {
Expand Down Expand Up @@ -424,7 +424,7 @@ func (s *panicSeekable) Size(_ context.Context) (int64, error) {
return int64(len(s.data)), nil
}

func (s *panicSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FrameTable, [32]byte, error) {
func (s *panicSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FullFrameTable, [32]byte, error) {
panic("not used")
}

Expand Down
3 changes: 2 additions & 1 deletion packages/orchestrator/pkg/sandbox/build_upload_v4.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,13 +75,14 @@ func (u *Upload) uploadFramed(
var selfBuild headers.BuildData

if srcPath != "" {
ft, checksum, err := storage.UploadFramed(ctx, u.store, u.paths.DataFile(string(fileType), cfg.CompressionType()), seekableTypeFor(fileType), srcPath, storage.WithCompressConfig(cfg), storage.WithMetadata(u.objectMetadata), storage.WithChecksumSHA256())
fullFT, checksum, err := storage.UploadFramed(ctx, u.store, u.paths.DataFile(string(fileType), cfg.CompressionType()), seekableTypeFor(fileType), srcPath, storage.WithCompressConfig(cfg), storage.WithMetadata(u.objectMetadata), storage.WithChecksumSHA256())
if err != nil {
return fmt.Errorf("%s upload: %w", fileType, err)
}

// Compressed: frame-table byte count, since sparse memfile diffs stream
// fewer bytes than they occupy on disk. Uncompressed has no table.
ft := fullFT.Table()
size := ft.UncompressedSize()
compressedSize := ft.CompressedSize()
if !ft.IsCompressed() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,7 @@ func (s *peerSeekable) OpenRangeReader(ctx context.Context, off int64, length in
return rc, err
}

func (s *peerSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FrameTable, [32]byte, error) {
func (s *peerSeekable) StoreFile(context.Context, string, ...storage.PutOption) (*storage.FullFrameTable, [32]byte, error) {
// peerSeekable only exists when routingProvider routed this buildID to an
// active peer at open time, i.e. the file is being P2P-served (the peer
// owns the upload). Asking the local orchestrator to upload it is a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ func TestPeerStorageProvider_FullTransitionFlow(t *testing.T) {

// 1. Pre-transition read via peer. ft={ct=None} (V3 header).
rc, err := seekable.OpenRangeReader(t.Context(), 0, int64(len(prePeerBytes)),
storage.NewFrameTable(storage.CompressionNone, nil))
storage.NewFullFrameTable(storage.CompressionNone, nil).Table())
require.NoError(t, err)
got, err := io.ReadAll(rc)
require.NoError(t, err)
Expand All @@ -210,14 +210,14 @@ func TestPeerStorageProvider_FullTransitionFlow(t *testing.T) {

// 2. First post-transition call: retriable error, no peer/base contact.
_, err = seekable.OpenRangeReader(t.Context(), 0, 1,
storage.NewFrameTable(storage.CompressionNone, nil))
storage.NewFullFrameTable(storage.CompressionNone, nil).Table())
var transErr *storage.PeerTransitionedError
require.ErrorAs(t, err, &transErr)

// 3. Caller reloads V4 header and retries with ct=Zstd. This must hit the
// compressed path on base.
rc, err = seekable.OpenRangeReader(t.Context(), 0, int64(len(postBaseBytes)),
storage.NewFrameTable(storage.CompressionZstd, nil))
storage.NewFullFrameTable(storage.CompressionZstd, nil).Table())
require.NoError(t, err)
got, err = io.ReadAll(rc)
require.NoError(t, err)
Expand Down
2 changes: 1 addition & 1 deletion packages/shared/pkg/storage/compress_encode.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ func newCompressorPool(cfg CompressConfig) (*sync.Pool, error) {
return pool, nil
}

func CompressBytes(ctx context.Context, data []byte, cfg CompressConfig) (*FrameTable, []byte, [32]byte, error) {
func CompressBytes(ctx context.Context, data []byte, cfg CompressConfig) (*FullFrameTable, []byte, [32]byte, error) {
up := &memPartUploader{}

const compressBytesConcurrency = 1
Expand Down
25 changes: 20 additions & 5 deletions packages/shared/pkg/storage/compress_frame_table.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,16 +70,31 @@ type FrameTable struct {
entries []frameEntry // sorted by StartU
}

// FullFrameTable marks a FrameTable that covers an entire file with no gaps.
// Produced by compressStream / CompressBytes / StoreFile / UploadFramed.
// The inner FrameTable is unexported and unembedded so methods are not
// promoted; consumers must reach functionality through nil-safe Table().
type FullFrameTable struct{ ft FrameTable }

// Table returns the underlying *FrameTable, nil-safely.
func (ft *FullFrameTable) Table() *FrameTable {
if ft == nil {
return nil
}

return &ft.ft
}

// newFrameTableFromEntries creates a FrameTable from pre-computed absolute-offset entries.
func newFrameTableFromEntries(ct CompressionType, entries []frameEntry) *FrameTable {
return &FrameTable{compressionType: ct, entries: entries}
}

// NewFrameTable creates a FrameTable from consecutive frame sizes, computing
// absolute offsets starting from zero.
func NewFrameTable(ct CompressionType, sizes []FrameSize) *FrameTable {
// NewFullFrameTable creates a FullFrameTable from consecutive frame sizes,
// computing absolute offsets starting from zero.
func NewFullFrameTable(ct CompressionType, sizes []FrameSize) *FullFrameTable {
if len(sizes) == 0 {
return newFrameTableFromEntries(ct, nil)
return &FullFrameTable{ft: *newFrameTableFromEntries(ct, nil)}
}

entries := make([]frameEntry, len(sizes))
Expand All @@ -96,7 +111,7 @@ func NewFrameTable(ct CompressionType, sizes []FrameSize) *FrameTable {
c += int64(s.C)
}

return newFrameTableFromEntries(ct, entries)
return &FullFrameTable{ft: *newFrameTableFromEntries(ct, entries)}
}

// CompressionType returns the compression type. Nil-safe: returns CompressionNone for nil.
Expand Down
14 changes: 7 additions & 7 deletions packages/shared/pkg/storage/compress_frame_table_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,10 +93,10 @@ func TestLocate(t *testing.T) {
func TestNewFrameTable(t *testing.T) {
t.Parallel()

ft := NewFrameTable(CompressionZstd, []FrameSize{
ft := NewFullFrameTable(CompressionZstd, []FrameSize{
{U: 1 << 20, C: 500_000},
{U: 1 << 20, C: 600_000},
})
}).Table()

require.Equal(t, 2, ft.NumFrames())
require.Equal(t, CompressionZstd, ft.CompressionType())
Expand All @@ -118,17 +118,17 @@ func TestNewFrameTable(t *testing.T) {
func TestFrameTable_TrimToRanges(t *testing.T) {
t.Parallel()

ft := NewFrameTable(CompressionLZ4, []FrameSize{
ft := NewFullFrameTable(CompressionLZ4, []FrameSize{
{U: 1 << 20, C: 500_000},
{U: 1 << 20, C: 600_000},
{U: 1 << 20, C: 400_000},
{U: 1 << 20, C: 700_000},
})
}).Table()

t.Run("all frames retained", func(t *testing.T) {
t.Parallel()
trimmed := ft.TrimToRanges([]Range{{Offset: 0, Length: 4 << 20}})
require.Same(t, ft, trimmed)
require.Equal(t, ft.NumFrames(), trimmed.NumFrames())
})

t.Run("single range trims to subset", func(t *testing.T) {
Expand Down Expand Up @@ -190,10 +190,10 @@ func TestSerializeDeserializeFrameTable(t *testing.T) {

t.Run("round-trip", func(t *testing.T) {
t.Parallel()
ft := NewFrameTable(CompressionZstd, []FrameSize{
ft := NewFullFrameTable(CompressionZstd, []FrameSize{
{U: 2048, C: 1024},
{U: 4096, C: 3500},
})
}).Table()

var buf bytes.Buffer
require.NoError(t, ft.Serialize(&buf))
Expand Down
4 changes: 2 additions & 2 deletions packages/shared/pkg/storage/compress_upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ func (p *part) addFrame(ctx context.Context, buf inputBuf, n int, pool *sync.Poo
})
}

func compressStream(ctx context.Context, in io.Reader, cfg CompressConfig, uploader partUploader, maxUploadConcurrency int, sink FrameSink) (*FrameTable, [32]byte, error) {
func compressStream(ctx context.Context, in io.Reader, cfg CompressConfig, uploader partUploader, maxUploadConcurrency int, sink FrameSink) (*FullFrameTable, [32]byte, error) {
ctx, cancel := context.WithCancel(ctx)
defer cancel()

Expand Down Expand Up @@ -179,7 +179,7 @@ func compressStream(ctx context.Context, in io.Reader, cfg CompressConfig, uploa
return nil, [32]byte{}, fmt.Errorf("complete upload: %w", err)
}

ft := NewFrameTable(cfg.CompressionType(), frameSizes)
ft := NewFullFrameTable(cfg.CompressionType(), frameSizes)

return ft, sum256(hasher), nil
}
Expand Down
8 changes: 5 additions & 3 deletions packages/shared/pkg/storage/compress_upload_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ func TestCompressStreamRoundTrip(t *testing.T) {
up := &memPartUploader{}
cfg := defaultCfg(tc.codec, tc.workers, tc.frameSize)

ft, checksum, err := compressStream(
fullFT, checksum, err := compressStream(
t.Context(),
bytes.NewReader(original),
cfg,
Expand All @@ -169,6 +169,7 @@ func TestCompressStreamRoundTrip(t *testing.T) {
nil,
)
require.NoError(t, err)
ft := fullFT.Table()

if tc.dataSize == 0 {
require.Equal(t, 0, ft.NumFrames())
Expand Down Expand Up @@ -300,7 +301,7 @@ func TestCompressStreamRace(t *testing.T) {
return fmt.Errorf("stream %d: checksum mismatch", i)
}

decompressed, err := decompressAll(ft, up.Assemble())
decompressed, err := decompressAll(ft.Table(), up.Assemble())
if err != nil {
return fmt.Errorf("stream %d: decompress: %w", i, err)
}
Expand Down Expand Up @@ -417,10 +418,11 @@ func BenchmarkStoreFile(b *testing.B) {
outPath := filepath.Join(outDir, "output.dat")
obj := &fsObject{path: outPath}

ft, _, err := obj.StoreFile(b.Context(), inputPath, WithCompressConfig(compCfg))
fullFT, _, err := obj.StoreFile(b.Context(), inputPath, WithCompressConfig(compCfg))
if err != nil {
b.Fatal(err)
}
ft := fullFT.Table()

b.ReportMetric(float64(ft.CompressedSize())/float64(ft.UncompressedSize()), "ratio")
}
Expand Down
Loading
Loading