From b4e6327911847e3897cfe7cb9ec5ab83c5c5e830 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:22:07 -0700 Subject: [PATCH 1/6] feat(orchestrator): record upload compression metrics --- .../orchestrator/pkg/sandbox/build_upload.go | 2 + .../pkg/sandbox/build_upload_v3.go | 25 +++++-- .../pkg/sandbox/build_upload_v4.go | 5 +- .../pkg/sandbox/upload_metrics.go | 66 +++++++++++++++++++ packages/shared/pkg/telemetry/meters.go | 12 ++++ 5 files changed, 105 insertions(+), 5 deletions(-) create mode 100644 packages/orchestrator/pkg/sandbox/upload_metrics.go diff --git a/packages/orchestrator/pkg/sandbox/build_upload.go b/packages/orchestrator/pkg/sandbox/build_upload.go index b481605747..ae7f295d52 100644 --- a/packages/orchestrator/pkg/sandbox/build_upload.go +++ b/packages/orchestrator/pkg/sandbox/build_upload.go @@ -25,6 +25,7 @@ type Upload struct { store storage.StorageProvider mem storage.CompressConfig root storage.CompressConfig + useCase string objectMetadata storage.ObjectMetadata future *utils.ErrorOnce useV4 bool @@ -57,6 +58,7 @@ func NewUpload( store: store, mem: mem, root: root, + useCase: useCase, objectMetadata: objectMetadata, useV4: memV4 || rootV4, } diff --git a/packages/orchestrator/pkg/sandbox/build_upload_v3.go b/packages/orchestrator/pkg/sandbox/build_upload_v3.go index 608be2bafa..56c6679ca9 100644 --- a/packages/orchestrator/pkg/sandbox/build_upload_v3.go +++ b/packages/orchestrator/pkg/sandbox/build_upload_v3.go @@ -5,6 +5,7 @@ package sandbox import ( "context" "fmt" + "os" "golang.org/x/sync/errgroup" @@ -31,7 +32,7 @@ func (u *Upload) runV3(ctx context.Context) error { return nil } - return headers.StoreHeader(egCtx, u.store, u.paths.MemfileHeader(), finalizeV3(u.snap.MemfileDiffHeader)) + return storeHeaderWithMetrics(egCtx, u.store, u.paths.MemfileHeader(), string(build.Memfile), u.useCase, finalizeV3(u.snap.MemfileDiffHeader)) }) eg.Go(func() error { @@ -39,7 +40,7 @@ func (u *Upload) runV3(ctx context.Context) error { return nil } - return headers.StoreHeader(egCtx, u.store, u.paths.RootfsHeader(), finalizeV3(u.snap.RootfsDiffHeader)) + return storeHeaderWithMetrics(egCtx, u.store, u.paths.RootfsHeader(), string(build.Rootfs), u.useCase, finalizeV3(u.snap.RootfsDiffHeader)) }) meta := storage.WithMetadata(u.objectMetadata) @@ -50,8 +51,16 @@ func (u *Upload) runV3(ctx context.Context) error { } _, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Memfile(), storage.MemfileObjectType, memfilePath, meta) + if err != nil { + return err + } + info, err := os.Stat(memfilePath) + if err != nil { + return fmt.Errorf("memfile stat: %w", err) + } + recordUploadCompression(egCtx, uploadArtifactData, string(build.Memfile), u.useCase, storage.CompressConfig{}, info.Size(), info.Size()) - return err + return nil }) eg.Go(func() error { @@ -60,8 +69,16 @@ func (u *Upload) runV3(ctx context.Context) error { } _, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Rootfs(), storage.RootFSObjectType, rootfsPath, meta) + if err != nil { + return err + } + info, err := os.Stat(rootfsPath) + if err != nil { + return fmt.Errorf("rootfs stat: %w", err) + } + recordUploadCompression(egCtx, uploadArtifactData, string(build.Rootfs), u.useCase, storage.CompressConfig{}, info.Size(), info.Size()) - return err + return nil }) eg.Go(func() error { diff --git a/packages/orchestrator/pkg/sandbox/build_upload_v4.go b/packages/orchestrator/pkg/sandbox/build_upload_v4.go index 7f5714d6dc..552f4ec4ac 100644 --- a/packages/orchestrator/pkg/sandbox/build_upload_v4.go +++ b/packages/orchestrator/pkg/sandbox/build_upload_v4.go @@ -71,14 +71,17 @@ func (u *Upload) uploadFramed( // Compressed: frame-table byte count, since sparse memfile diffs stream // fewer bytes than they occupy on disk. Uncompressed has no table. size := ft.UncompressedSize() + compressedSize := ft.CompressedSize() if !ft.IsCompressed() { info, statErr := os.Stat(srcPath) if statErr != nil { return fmt.Errorf("%s stat: %w", fileType, statErr) } size = info.Size() + compressedSize = size } + recordUploadCompression(ctx, uploadArtifactData, string(fileType), u.useCase, cfg, size, compressedSize) selfBuild = headers.BuildData{Size: size, Checksum: checksum, FrameData: ft} } @@ -93,7 +96,7 @@ func (u *Upload) uploadFramed( } h.Builds[u.buildID] = selfBuild - if err := headers.StoreHeader(ctx, u.store, u.paths.HeaderFile(string(fileType)), h); err != nil { + if err := storeHeaderWithMetrics(ctx, u.store, u.paths.HeaderFile(string(fileType)), string(fileType), u.useCase, h); err != nil { return fmt.Errorf("store %s header: %w", fileType, err) } diff --git a/packages/orchestrator/pkg/sandbox/upload_metrics.go b/packages/orchestrator/pkg/sandbox/upload_metrics.go new file mode 100644 index 0000000000..261092283c --- /dev/null +++ b/packages/orchestrator/pkg/sandbox/upload_metrics.go @@ -0,0 +1,66 @@ +//go:build linux + +package sandbox + +import ( + "context" + "fmt" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" + + "github.com/e2b-dev/infra/packages/shared/pkg/storage" + headers "github.com/e2b-dev/infra/packages/shared/pkg/storage/header" + "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" + "github.com/e2b-dev/infra/packages/shared/pkg/utils" +) + +const ( + uploadArtifactData = "data" + uploadArtifactHeader = "header" +) + +var ( + uploadUncompressedBytes = utils.Must(telemetry.GetHistogram(meter, telemetry.UploadUncompressedBytes)) + uploadCompressedBytes = utils.Must(telemetry.GetHistogram(meter, telemetry.UploadCompressedBytes)) + uploadCompressionRatioBp = utils.Must(telemetry.GetHistogram(meter, telemetry.UploadCompressionRatioBp)) +) + +func recordUploadCompression(ctx context.Context, artifact, fileType, useCase string, cfg storage.CompressConfig, uncompressed, compressed int64) { + attrs := metric.WithAttributes( + attribute.String("artifact", artifact), + attribute.String("file_type", fileType), + attribute.String("use_case", useCase), + attribute.String("compression.type", cfg.CompressionType().String()), + attribute.Int("compression.level", cfg.Level), + ) + + uploadUncompressedBytes.Record(ctx, uncompressed, attrs) + uploadCompressedBytes.Record(ctx, compressed, attrs) + uploadCompressionRatioBp.Record(ctx, ratioBp(compressed, uncompressed), attrs) +} + +func storeHeaderWithMetrics(ctx context.Context, store storage.StorageProvider, path, fileType, useCase string, h *headers.Header) error { + if h.IncompletePendingUpload { + return fmt.Errorf("refusing to persist incomplete header for %s", path) + } + + data, err := headers.SerializeHeader(h) + if err != nil { + return fmt.Errorf("serialize header: %w", err) + } + + blob, err := store.OpenBlob(ctx, path, storage.MetadataObjectType) + if err != nil { + return fmt.Errorf("open blob %s: %w", path, err) + } + + if err := blob.Put(ctx, data); err != nil { + return err + } + + size := int64(len(data)) + recordUploadCompression(ctx, uploadArtifactHeader, fileType, useCase, storage.CompressConfig{}, size, size) + + return nil +} diff --git a/packages/shared/pkg/telemetry/meters.go b/packages/shared/pkg/telemetry/meters.go index c6a8c0c46a..9b35eec015 100644 --- a/packages/shared/pkg/telemetry/meters.go +++ b/packages/shared/pkg/telemetry/meters.go @@ -125,6 +125,10 @@ const ( SnapshotDiffBytes HistogramType = "orchestrator.sandbox.snapshot.diff.bytes" SnapshotDiffRatioBp HistogramType = "orchestrator.sandbox.snapshot.diff.ratio_bp" SnapshotTotalBytes HistogramType = "orchestrator.sandbox.snapshot.total.bytes" + + UploadUncompressedBytes HistogramType = "orchestrator.sandbox.upload.uncompressed.bytes" + UploadCompressedBytes HistogramType = "orchestrator.sandbox.upload.compressed.bytes" + UploadCompressionRatioBp HistogramType = "orchestrator.sandbox.upload.compression.ratio_bp" ) const ( @@ -363,6 +367,10 @@ var histogramDesc = map[HistogramType]string{ SnapshotDiffBytes: "Per-snapshot dirty/empty bytes per file", SnapshotDiffRatioBp: "Per-snapshot dirty/empty as fraction of total mapped size, in basis points (10000=100%)", SnapshotTotalBytes: "Per-snapshot total mapped size of the file", + + UploadUncompressedBytes: "Per-upload uncompressed artifact size", + UploadCompressedBytes: "Per-upload compressed artifact size", + UploadCompressionRatioBp: "Per-upload compressed/uncompressed ratio, in basis points (10000=100%)", } var histogramUnits = map[HistogramType]string{ @@ -397,6 +405,10 @@ var histogramUnits = map[HistogramType]string{ SnapshotDiffBytes: "{By}", SnapshotDiffRatioBp: "{1}", SnapshotTotalBytes: "{By}", + + UploadUncompressedBytes: "{By}", + UploadCompressedBytes: "{By}", + UploadCompressionRatioBp: "{1}", } func GetHistogram(meter metric.Meter, name HistogramType) (metric.Int64Histogram, error) { From 4dba885213af2eae821891c07bf49b75ec6d30e1 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:33:13 -0700 Subject: [PATCH 2/6] refactor(orchestrator): simplify header upload metrics --- .../orchestrator/pkg/sandbox/upload_metrics.go | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) diff --git a/packages/orchestrator/pkg/sandbox/upload_metrics.go b/packages/orchestrator/pkg/sandbox/upload_metrics.go index 261092283c..72ebeea143 100644 --- a/packages/orchestrator/pkg/sandbox/upload_metrics.go +++ b/packages/orchestrator/pkg/sandbox/upload_metrics.go @@ -4,7 +4,6 @@ package sandbox import ( "context" - "fmt" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" @@ -41,21 +40,12 @@ func recordUploadCompression(ctx context.Context, artifact, fileType, useCase st } func storeHeaderWithMetrics(ctx context.Context, store storage.StorageProvider, path, fileType, useCase string, h *headers.Header) error { - if h.IncompletePendingUpload { - return fmt.Errorf("refusing to persist incomplete header for %s", path) + if err := headers.StoreHeader(ctx, store, path, h); err != nil { + return err } data, err := headers.SerializeHeader(h) if err != nil { - return fmt.Errorf("serialize header: %w", err) - } - - blob, err := store.OpenBlob(ctx, path, storage.MetadataObjectType) - if err != nil { - return fmt.Errorf("open blob %s: %w", path, err) - } - - if err := blob.Put(ctx, data); err != nil { return err } From 0d317160488dfeff8e315cafcce85a3bfff3198c Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:36:01 -0700 Subject: [PATCH 3/6] fix(orchestrator): preserve upload expansion ratios --- packages/orchestrator/pkg/sandbox/upload_metrics.go | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/packages/orchestrator/pkg/sandbox/upload_metrics.go b/packages/orchestrator/pkg/sandbox/upload_metrics.go index 72ebeea143..114669dc88 100644 --- a/packages/orchestrator/pkg/sandbox/upload_metrics.go +++ b/packages/orchestrator/pkg/sandbox/upload_metrics.go @@ -36,7 +36,15 @@ func recordUploadCompression(ctx context.Context, artifact, fileType, useCase st uploadUncompressedBytes.Record(ctx, uncompressed, attrs) uploadCompressedBytes.Record(ctx, compressed, attrs) - uploadCompressionRatioBp.Record(ctx, ratioBp(compressed, uncompressed), attrs) + uploadCompressionRatioBp.Record(ctx, uploadRatioBp(compressed, uncompressed), attrs) +} + +func uploadRatioBp(compressed, uncompressed int64) int64 { + if uncompressed <= 0 || compressed < 0 { + return 0 + } + + return compressed * 10000 / uncompressed } func storeHeaderWithMetrics(ctx context.Context, store storage.StorageProvider, path, fileType, useCase string, h *headers.Header) error { From d66ee67369ac94f1dfff86b37a6183922797c30e Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:40:28 -0700 Subject: [PATCH 4/6] fix(telemetry): bucket upload byte histograms --- packages/shared/pkg/telemetry/metrics.go | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/packages/shared/pkg/telemetry/metrics.go b/packages/shared/pkg/telemetry/metrics.go index b66a6a8bb4..6779790748 100644 --- a/packages/shared/pkg/telemetry/metrics.go +++ b/packages/shared/pkg/telemetry/metrics.go @@ -56,7 +56,7 @@ func NewMeterExporter(ctx context.Context, extraOption ...otlpmetricgrpc.Option) // snapshotBytesView routes the per-snapshot byte histograms to a base-2 // exponential aggregation. The SDK's default explicit buckets max out at // 10_000 (tuned for ms), which collapses byte values into +Inf and makes -// percentile queries unusable for the 0–tens-of-GiB range these metrics +// percentile queries unusable for the large byte ranges these metrics // cover. var snapshotBytesView = sdkmetric.NewView( sdkmetric.Instrument{ @@ -72,6 +72,20 @@ var snapshotBytesView = sdkmetric.NewView( }, ) +var uploadBytesView = sdkmetric.NewView( + sdkmetric.Instrument{ + Kind: sdkmetric.InstrumentKindHistogram, + Name: "orchestrator.sandbox.upload.*.bytes", + Unit: "{By}", + }, + sdkmetric.Stream{ + Aggregation: sdkmetric.AggregationBase2ExponentialHistogram{ + MaxSize: 160, + MaxScale: 20, + }, + }, +) + func NewMeterProvider(metricsExporter sdkmetric.Exporter, metricExportPeriod time.Duration, res *resource.Resource, extraOption ...sdkmetric.Option) (metric.MeterProvider, error) { opts := []sdkmetric.Option{ sdkmetric.WithReader( @@ -91,7 +105,7 @@ func NewMeterProvider(metricsExporter sdkmetric.Exporter, metricExportPeriod tim } opts = append(opts, extraOption...) - opts = append(opts, sdkmetric.WithView(snapshotBytesView)) + opts = append(opts, sdkmetric.WithView(snapshotBytesView, uploadBytesView)) return sdkmetric.NewMeterProvider(opts...), nil } From 7cf2094eb6e464e5a3bf3f1c7184e8c82102560e Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:45:44 -0700 Subject: [PATCH 5/6] fix(telemetry): match upload byte histogram view by unit --- packages/shared/pkg/telemetry/metrics.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/shared/pkg/telemetry/metrics.go b/packages/shared/pkg/telemetry/metrics.go index 6779790748..76a88668a9 100644 --- a/packages/shared/pkg/telemetry/metrics.go +++ b/packages/shared/pkg/telemetry/metrics.go @@ -75,7 +75,7 @@ var snapshotBytesView = sdkmetric.NewView( var uploadBytesView = sdkmetric.NewView( sdkmetric.Instrument{ Kind: sdkmetric.InstrumentKindHistogram, - Name: "orchestrator.sandbox.upload.*.bytes", + Name: "orchestrator.sandbox.upload.*", Unit: "{By}", }, sdkmetric.Stream{ From c1242ad7283a5374db319808fe77a3047986b6cd Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Wed, 20 May 2026 12:48:51 -0700 Subject: [PATCH 6/6] fix(orchestrator): stat v3 upload inputs before upload --- .../orchestrator/pkg/sandbox/build_upload_v3.go | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/packages/orchestrator/pkg/sandbox/build_upload_v3.go b/packages/orchestrator/pkg/sandbox/build_upload_v3.go index 56c6679ca9..0ea56aeb44 100644 --- a/packages/orchestrator/pkg/sandbox/build_upload_v3.go +++ b/packages/orchestrator/pkg/sandbox/build_upload_v3.go @@ -50,14 +50,14 @@ func (u *Upload) runV3(ctx context.Context) error { return nil } - _, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Memfile(), storage.MemfileObjectType, memfilePath, meta) - if err != nil { - return err - } info, err := os.Stat(memfilePath) if err != nil { return fmt.Errorf("memfile stat: %w", err) } + _, _, err = storage.UploadFramed(egCtx, u.store, u.paths.Memfile(), storage.MemfileObjectType, memfilePath, meta) + if err != nil { + return err + } recordUploadCompression(egCtx, uploadArtifactData, string(build.Memfile), u.useCase, storage.CompressConfig{}, info.Size(), info.Size()) return nil @@ -68,14 +68,14 @@ func (u *Upload) runV3(ctx context.Context) error { return nil } - _, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Rootfs(), storage.RootFSObjectType, rootfsPath, meta) - if err != nil { - return err - } info, err := os.Stat(rootfsPath) if err != nil { return fmt.Errorf("rootfs stat: %w", err) } + _, _, err = storage.UploadFramed(egCtx, u.store, u.paths.Rootfs(), storage.RootFSObjectType, rootfsPath, meta) + if err != nil { + return err + } recordUploadCompression(egCtx, uploadArtifactData, string(build.Rootfs), u.useCase, storage.CompressConfig{}, info.Size(), info.Size()) return nil