diff --git a/packages/orchestrator/pkg/sandbox/upload_metrics.go b/packages/orchestrator/pkg/sandbox/upload_metrics.go index 499a85200e..6041a5f9f9 100644 --- a/packages/orchestrator/pkg/sandbox/upload_metrics.go +++ b/packages/orchestrator/pkg/sandbox/upload_metrics.go @@ -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 } diff --git a/packages/shared/pkg/storage/header/serialization.go b/packages/shared/pkg/storage/header/serialization.go index 9e25212495..cf0017d92b 100644 --- a/packages/shared/pkg/storage/header/serialization.go +++ b/packages/shared/pkg/storage/header/serialization.go @@ -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. @@ -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). diff --git a/packages/shared/pkg/storage/header/serialization_v4.go b/packages/shared/pkg/storage/header/serialization_v4.go index a88a772ede..e89dab83e6 100644 --- a/packages/shared/pkg/storage/header/serialization_v4.go +++ b/packages/shared/pkg/storage/header/serialization_v4.go @@ -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 @@ -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) @@ -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 { @@ -97,7 +98,7 @@ 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) } } @@ -105,7 +106,7 @@ func serializeV4(metadata *Metadata, builds map[uuid.UUID]BuildData, mappings [] 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 @@ -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.