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
4 changes: 2 additions & 2 deletions packages/orchestrator/pkg/sandbox/upload_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,12 +47,12 @@ func uploadRatioBp(compressed, uncompressed int64) int64 {
}

func storeHeaderWithMetrics(ctx context.Context, store storage.StorageProvider, path, fileType string, h *headers.Header) error {
size, err := headers.StoreHeader(ctx, store, path, h)
cfg, stored, uncompressed, err := headers.StoreHeader(ctx, store, path, h)
if err != nil {
return err
}

recordUploadCompression(ctx, uploadArtifactHeader, fileType, storage.CompressConfig{}, size, size)
recordUploadCompression(ctx, uploadArtifactHeader, fileType, cfg, uncompressed, stored)

return nil
}
41 changes: 28 additions & 13 deletions packages/shared/pkg/storage/header/serialization.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@ func SerializeHeader(h *Header) ([]byte, error) {
return serializeV3(h.Metadata, h.Mapping)
}

return serializeV4(h.Metadata, h.Builds, h.Mapping, h.IncompletePendingUpload)
data, _, err := serializeV4(h.Metadata, h.Builds, h.Mapping, h.IncompletePendingUpload)

return data, err
}

// DeserializeBytes auto-detects the header version and deserializes accordingly.
Expand Down Expand Up @@ -57,33 +59,46 @@ func LoadHeader(ctx context.Context, s storage.StorageProvider, path string) (*H
return DeserializeBytes(data)
}

// StoreHeader serializes a header, uploads it, and returns the stored byte count.
// Refuses to persist a header still flagged as in-flight — the upload pipeline
// must clear IncompletePendingUpload before reaching here.
func StoreHeader(ctx context.Context, s storage.StorageProvider, path string, h *Header) (int64, error) {
// StoreHeader serializes a header, uploads it, and returns the effective
// compression config plus the stored and pre-compression byte counts. V3 has
// no inner compression so the counts match and cfg is the zero value. Refuses
// to persist a header still flagged as in-flight.
func StoreHeader(ctx context.Context, s storage.StorageProvider, path string, h *Header) (cfg storage.CompressConfig, stored, uncompressed int64, err error) {
if h == nil {
return 0, errors.New("header is nil")
return storage.CompressConfig{}, 0, 0, errors.New("header is nil")
}

if h.IncompletePendingUpload {
return 0, fmt.Errorf("refusing to persist incomplete header for %s", path)
return storage.CompressConfig{}, 0, 0, fmt.Errorf("refusing to persist incomplete header for %s", path)
}

data, err := SerializeHeader(h)
if err != nil {
return 0, fmt.Errorf("serialize header: %w", err)
var data []byte
if h.Metadata.Version <= 3 {
data, err = serializeV3(h.Metadata, h.Mapping)
if err != nil {
return storage.CompressConfig{}, 0, 0, fmt.Errorf("serialize header: %w", err)
}
uncompressed = int64(len(data))
} else {
var blockUncompressed int64
data, blockUncompressed, err = serializeV4(h.Metadata, h.Builds, h.Mapping, h.IncompletePendingUpload)
if err != nil {
return storage.CompressConfig{}, 0, 0, fmt.Errorf("serialize header: %w", err)
}
uncompressed = int64(metadataSize+v4FlagsLen+v4SizePrefixLen) + blockUncompressed
cfg.Type = storage.CompressionLZ4.String()
}

blob, err := s.OpenBlob(ctx, path, storage.MetadataObjectType)
if err != nil {
return 0, fmt.Errorf("open blob %s: %w", path, err)
return storage.CompressConfig{}, 0, 0, fmt.Errorf("open blob %s: %w", path, err)
}

if err := blob.Put(ctx, data); err != nil {
return 0, fmt.Errorf("put blob %s: %w", path, err)
return storage.CompressConfig{}, 0, 0, fmt.Errorf("put blob %s: %w", path, err)
}

return int64(len(data)), nil
return cfg, int64(len(data)), uncompressed, nil
}

// Deserialize reads a header from a storage Blob (legacy API).
Expand Down
19 changes: 10 additions & 9 deletions packages/shared/pkg/storage/header/serialization_v4.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,10 +43,11 @@ type v4SerializableBuildInfo struct {

// serializeV4 writes [Metadata] [uint8 flags] [uint32 LZ4 size] [LZ4( Builds[] + Mappings[] )].
// Frame tables are sparse-trimmed to only frames referenced by mappings.
func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []BuildMap, incomplete bool) ([]byte, error) {
// Also returns the uncompressed inner-block size (LZ4 input length).
func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []BuildMap, incomplete bool) ([]byte, int64, error) {
var metaBuf bytes.Buffer
if err := binary.Write(&metaBuf, binary.LittleEndian, metadata); err != nil {
return nil, fmt.Errorf("failed to write metadata: %w", err)
return nil, 0, fmt.Errorf("failed to write metadata: %w", err)
}

var block bytes.Buffer
Expand All @@ -61,7 +62,7 @@ func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []
})

if err := binary.Write(&block, binary.LittleEndian, uint32(len(buildIDs))); err != nil {
return nil, fmt.Errorf("failed to write build count: %w", err)
return nil, 0, fmt.Errorf("failed to write build count: %w", err)
}

buildRanges := extractRelevantRanges(mappings)
Expand All @@ -75,17 +76,17 @@ func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []
}

if err := binary.Write(&block, binary.LittleEndian, &entry); err != nil {
return nil, fmt.Errorf("failed to write build info: %w", err)
return nil, 0, fmt.Errorf("failed to write build info: %w", err)
}

trimmed := bd.FrameData.TrimToRanges(buildRanges[id])
if err := trimmed.Serialize(&block); err != nil {
return nil, fmt.Errorf("failed to write build frame data: %w", err)
return nil, 0, fmt.Errorf("failed to write build frame data: %w", err)
}
}

if err := binary.Write(&block, binary.LittleEndian, uint32(len(mappings))); err != nil {
return nil, fmt.Errorf("failed to write mappings count: %w", err)
return nil, 0, fmt.Errorf("failed to write mappings count: %w", err)
}

for _, mapping := range mappings {
Expand All @@ -97,15 +98,15 @@ func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []
}

if err := binary.Write(&block, binary.LittleEndian, v4); err != nil {
return nil, fmt.Errorf("failed to write block mapping: %w", err)
return nil, 0, fmt.Errorf("failed to write block mapping: %w", err)
}
}

// LZ4-compress the block and assemble: [metadata] [uint8 flags] [uint32 size] [compressed block].
blockBytes := block.Bytes()
compressed, err := compressLZ4(blockBytes)
if err != nil {
return nil, fmt.Errorf("failed to LZ4-compress v4 header block: %w", err)
return nil, 0, fmt.Errorf("failed to LZ4-compress v4 header block: %w", err)
}

var flags uint8
Expand All @@ -119,7 +120,7 @@ func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings []
binary.LittleEndian.PutUint32(result[metadataSize+v4FlagsLen:], uint32(len(blockBytes)))
copy(result[metadataSize+v4FlagsLen+v4SizePrefixLen:], compressed)

return result, nil
return result, int64(len(blockBytes)), nil
}

// deserializeV4 decompresses and reads the V4 block.
Expand Down
Loading