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
2 changes: 2 additions & 0 deletions packages/orchestrator/pkg/sandbox/build_upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -57,6 +58,7 @@ func NewUpload(
store: store,
mem: mem,
root: root,
useCase: useCase,
objectMetadata: objectMetadata,
useV4: memV4 || rootV4,
}
Expand Down
29 changes: 23 additions & 6 deletions packages/orchestrator/pkg/sandbox/build_upload_v3.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ package sandbox
import (
"context"
"fmt"
"os"

"golang.org/x/sync/errgroup"

Expand All @@ -31,15 +32,15 @@ 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 {
if u.snap.RootfsDiffHeader == nil {
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)
Expand All @@ -49,19 +50,35 @@ func (u *Upload) runV3(ctx context.Context) error {
return nil
}

_, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Memfile(), storage.MemfileObjectType, memfilePath, meta)
info, err := os.Stat(memfilePath)
if err != nil {
return fmt.Errorf("memfile stat: %w", err)
}
Comment thread
ValentaTomas marked this conversation as resolved.
_, _, 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 err
return nil
})

eg.Go(func() error {
if rootfsPath == "" {
return nil
}

_, _, err := storage.UploadFramed(egCtx, u.store, u.paths.Rootfs(), storage.RootFSObjectType, rootfsPath, meta)
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 err
return nil
})

eg.Go(func() error {
Expand Down
5 changes: 4 additions & 1 deletion packages/orchestrator/pkg/sandbox/build_upload_v4.go
Original file line number Diff line number Diff line change
Expand Up @@ -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}
}

Expand All @@ -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)
}

Expand Down
64 changes: 64 additions & 0 deletions packages/orchestrator/pkg/sandbox/upload_metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
//go:build linux

package sandbox

import (
"context"

"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, 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 {
if err := headers.StoreHeader(ctx, store, path, h); err != nil {
return err
}

data, err := headers.SerializeHeader(h)
if err != nil {
return err
}

size := int64(len(data))
recordUploadCompression(ctx, uploadArtifactHeader, fileType, useCase, storage.CompressConfig{}, size, size)

return nil
}
12 changes: 12 additions & 0 deletions packages/shared/pkg/telemetry/meters.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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{
Expand Down Expand Up @@ -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) {
Expand Down
18 changes: 16 additions & 2 deletions packages/shared/pkg/telemetry/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand All @@ -72,6 +72,20 @@ var snapshotBytesView = sdkmetric.NewView(
},
)

var uploadBytesView = sdkmetric.NewView(
sdkmetric.Instrument{
Kind: sdkmetric.InstrumentKindHistogram,
Name: "orchestrator.sandbox.upload.*",
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(
Expand All @@ -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
}
Loading