diff --git a/packages/orchestrator/pkg/sandbox/block/chunk.go b/packages/orchestrator/pkg/sandbox/block/chunk.go index c3f29ea15f..ad2017d2aa 100644 --- a/packages/orchestrator/pkg/sandbox/block/chunk.go +++ b/packages/orchestrator/pkg/sandbox/block/chunk.go @@ -8,6 +8,7 @@ import ( "strconv" "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" "go.uber.org/zap" "golang.org/x/sync/errgroup" "golang.org/x/sync/singleflight" @@ -17,8 +18,67 @@ import ( "github.com/e2b-dev/infra/packages/shared/pkg/logger" "github.com/e2b-dev/infra/packages/shared/pkg/storage" "github.com/e2b-dev/infra/packages/shared/pkg/storage/header" + "github.com/e2b-dev/infra/packages/shared/pkg/telemetry" ) +const ( + pullType = "pull-type" + pullTypeLocal = "local" + pullTypeRemote = "remote" + + failureReason = "failure-reason" + + failureTypeLocalRead = "local-read" + failureTypeLocalReadAgain = "local-read-again" + failureTypeRemoteRead = "remote-read" + failureTypeCacheFetch = "cache-fetch" +) + +type precomputedAttrs struct { + successFromCache metric.MeasurementOption + successFromRemote metric.MeasurementOption + + failCacheRead metric.MeasurementOption + failRemoteFetch metric.MeasurementOption + failLocalReadAgain metric.MeasurementOption + + // RemoteReads timer (runFetch) + remoteSuccess metric.MeasurementOption + remoteFailure metric.MeasurementOption +} + +var chunkerAttrs = precomputedAttrs{ + successFromCache: telemetry.PrecomputeAttrs( + telemetry.Success, + attribute.String(pullType, pullTypeLocal)), + + successFromRemote: telemetry.PrecomputeAttrs( + telemetry.Success, + attribute.String(pullType, pullTypeRemote)), + + failCacheRead: telemetry.PrecomputeAttrs( + telemetry.Failure, + attribute.String(pullType, pullTypeLocal), + attribute.String(failureReason, failureTypeLocalRead)), + + failRemoteFetch: telemetry.PrecomputeAttrs( + telemetry.Failure, + attribute.String(pullType, pullTypeRemote), + attribute.String(failureReason, failureTypeCacheFetch)), + + failLocalReadAgain: telemetry.PrecomputeAttrs( + telemetry.Failure, + attribute.String(pullType, pullTypeLocal), + attribute.String(failureReason, failureTypeLocalReadAgain)), + + remoteSuccess: telemetry.PrecomputeAttrs( + telemetry.Success), + + remoteFailure: telemetry.PrecomputeAttrs( + telemetry.Failure, + attribute.String(failureReason, failureTypeRemoteRead)), +} + // Chunker is the interface satisfied by both FullFetchChunker and StreamingChunker. type Chunker interface { Slice(ctx context.Context, off, length int64) ([]byte, error) @@ -125,40 +185,32 @@ func (c *FullFetchChunker) Slice(ctx context.Context, off, length int64) ([]byte b, err := c.cache.Slice(off, length) if err == nil { - timer.Success(ctx, length, - attribute.String(pullType, pullTypeLocal)) + timer.RecordRaw(ctx, length, chunkerAttrs.successFromCache) return b, nil } if !errors.As(err, &BytesNotAvailableError{}) { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeLocal), - attribute.String(failureReason, failureTypeLocalRead)) + timer.RecordRaw(ctx, length, chunkerAttrs.failCacheRead) return nil, fmt.Errorf("failed read from cache at offset %d: %w", off, err) } chunkErr := c.fetchToCache(ctx, off, length) if chunkErr != nil { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeRemote), - attribute.String(failureReason, failureTypeCacheFetch)) + timer.RecordRaw(ctx, length, chunkerAttrs.failRemoteFetch) return nil, fmt.Errorf("failed to ensure data at %d-%d: %w", off, off+length, chunkErr) } b, cacheErr := c.cache.Slice(off, length) if cacheErr != nil { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeLocal), - attribute.String(failureReason, failureTypeLocalReadAgain)) + timer.RecordRaw(ctx, length, chunkerAttrs.failLocalReadAgain) return nil, fmt.Errorf("failed to read from cache after ensuring data at %d-%d: %w", off, off+length, cacheErr) } - timer.Success(ctx, length, - attribute.String(pullType, pullTypeRemote)) + timer.RecordRaw(ctx, length, chunkerAttrs.successFromRemote) return b, nil } @@ -210,24 +262,20 @@ func (c *FullFetchChunker) fetchToCache(ctx context.Context, off, length int64) readBytes, err := c.base.ReadAt(ctx, b, fetchOff) if err != nil { - fetchSW.Failure(ctx, int64(readBytes), - attribute.String(failureReason, failureTypeRemoteRead), - ) + fetchSW.RecordRaw(ctx, int64(readBytes), chunkerAttrs.remoteFailure) return nil, fmt.Errorf("failed to read chunk from base %d: %w", fetchOff, err) } if readBytes != len(b) { - fetchSW.Failure(ctx, int64(readBytes), - attribute.String(failureReason, failureTypeRemoteRead), - ) + fetchSW.RecordRaw(ctx, int64(readBytes), chunkerAttrs.remoteFailure) return nil, fmt.Errorf("failed to read chunk from base %d: expected %d bytes, got %d bytes", fetchOff, len(b), readBytes) } c.cache.setIsCached(fetchOff, int64(readBytes)) - fetchSW.Success(ctx, int64(readBytes)) + fetchSW.RecordRaw(ctx, int64(readBytes), chunkerAttrs.remoteSuccess) return nil, nil }) @@ -251,16 +299,3 @@ func (c *FullFetchChunker) Close() error { func (c *FullFetchChunker) FileSize() (int64, error) { return c.cache.FileSize() } - -const ( - pullType = "pull-type" - pullTypeLocal = "local" - pullTypeRemote = "remote" - - failureReason = "failure-reason" - - failureTypeLocalRead = "local-read" - failureTypeLocalReadAgain = "local-read-again" - failureTypeRemoteRead = "remote-read" - failureTypeCacheFetch = "cache-fetch" -) diff --git a/packages/orchestrator/pkg/sandbox/block/chunk_bench_test.go b/packages/orchestrator/pkg/sandbox/block/chunk_bench_test.go new file mode 100644 index 0000000000..93534b4d3b --- /dev/null +++ b/packages/orchestrator/pkg/sandbox/block/chunk_bench_test.go @@ -0,0 +1,60 @@ +package block + +import ( + "context" + "path/filepath" + "testing" + + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + + blockmetrics "github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block/metrics" +) + +const ( + cbBlockSize int64 = 4096 + cbNumBlocks int64 = 16384 // 64 MiB + cbCacheSize int64 = cbNumBlocks * cbBlockSize + cbChunkSize int64 = 4 * 1024 * 1024 // 4 MiB — MemoryChunkSize + cbChunkCount int64 = cbCacheSize / cbChunkSize +) + +// BenchmarkChunkerSlice_CacheHit benchmarks the full FullFetchChunker.Slice +// hot path on a cache hit: bitmap check + mmap slice return + OTEL +// timer.Success with attribute construction. +func BenchmarkChunkerSlice_CacheHit(b *testing.B) { + provider := sdkmetric.NewMeterProvider() + b.Cleanup(func() { provider.Shutdown(context.Background()) }) + + m, err := blockmetrics.NewMetrics(provider) + if err != nil { + b.Fatal(err) + } + + chunker, err := NewFullFetchChunker( + cbCacheSize, cbBlockSize, + nil, // base is never called on cache hit + filepath.Join(b.TempDir(), "cache"), + m, + ) + if err != nil { + b.Fatal(err) + } + b.Cleanup(func() { chunker.Close() }) + + // Pre-populate the cache so every Slice hits. + chunker.cache.setIsCached(0, cbCacheSize) + + ctx := context.Background() + + b.ResetTimer() + for i := range b.N { + off := int64(i%int(cbChunkCount)) * cbChunkSize + s, sliceErr := chunker.Slice(ctx, off, cbChunkSize) + if sliceErr != nil { + b.Fatal(sliceErr) + } + if len(s) == 0 { + b.Fatal("empty slice") + } + } +} diff --git a/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go b/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go index 956d71e0b3..7e40b35c4e 100644 --- a/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go +++ b/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go @@ -11,7 +11,6 @@ import ( "sync/atomic" "time" - "go.opentelemetry.io/otel/attribute" "golang.org/x/sync/errgroup" "github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block/metrics" @@ -227,16 +226,13 @@ func (c *StreamingChunker) Slice(ctx context.Context, off, length int64) ([]byte // Fast path: already cached b, err := c.cache.Slice(off, length) if err == nil { - timer.Success(ctx, length, - attribute.String(pullType, pullTypeLocal)) + timer.RecordRaw(ctx, length, chunkerAttrs.successFromCache) return b, nil } if !errors.As(err, &BytesNotAvailableError{}) { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeLocal), - attribute.String(failureReason, failureTypeLocalRead)) + timer.RecordRaw(ctx, length, chunkerAttrs.failCacheRead) return nil, fmt.Errorf("failed read from cache at offset %d: %w", off, err) } @@ -269,24 +265,19 @@ func (c *StreamingChunker) Slice(ctx context.Context, off, length int64) ([]byte } if err := eg.Wait(); err != nil { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeRemote), - attribute.String(failureReason, failureTypeCacheFetch)) + timer.RecordRaw(ctx, length, chunkerAttrs.failRemoteFetch) return nil, fmt.Errorf("failed to ensure data at %d-%d: %w", off, off+length, err) } b, cacheErr := c.cache.Slice(off, length) if cacheErr != nil { - timer.Failure(ctx, length, - attribute.String(pullType, pullTypeLocal), - attribute.String(failureReason, failureTypeLocalReadAgain)) + timer.RecordRaw(ctx, length, chunkerAttrs.failLocalReadAgain) return nil, fmt.Errorf("failed to read from cache after ensuring data at %d-%d: %w", off, off+length, cacheErr) } - timer.Success(ctx, length, - attribute.String(pullType, pullTypeRemote)) + timer.RecordRaw(ctx, length, chunkerAttrs.successFromRemote) return b, nil } @@ -386,15 +377,14 @@ func (c *StreamingChunker) runFetch(ctx context.Context, s *fetchSession) { err = c.progressiveRead(ctx, s, mmapSlice) if err != nil { - fetchTimer.Failure(ctx, s.chunkLen, - attribute.String(failureReason, failureTypeRemoteRead)) + fetchTimer.RecordRaw(ctx, s.chunkLen, chunkerAttrs.remoteFailure) s.setError(err, false) return } - fetchTimer.Success(ctx, s.chunkLen) + fetchTimer.RecordRaw(ctx, s.chunkLen, chunkerAttrs.remoteSuccess) s.setDone() } diff --git a/packages/shared/pkg/telemetry/meters.go b/packages/shared/pkg/telemetry/meters.go index 70ad5da013..b6174169c4 100644 --- a/packages/shared/pkg/telemetry/meters.go +++ b/packages/shared/pkg/telemetry/meters.go @@ -411,6 +411,12 @@ const ( resultTypeFailure = "failure" ) +var ( + // Pre-allocated result attributes for use with PrecomputeAttrs. + Success = attribute.String(resultAttr, resultTypeSuccess) + Failure = attribute.String(resultAttr, resultTypeFailure) +) + func (t Stopwatch) Success(ctx context.Context, total int64, kv ...attribute.KeyValue) { t.end(ctx, resultTypeSuccess, total, kv...) } @@ -422,9 +428,23 @@ func (t Stopwatch) Failure(ctx context.Context, total int64, kv ...attribute.Key func (t Stopwatch) end(ctx context.Context, result string, total int64, kv ...attribute.KeyValue) { kv = append(kv, attribute.KeyValue{Key: resultAttr, Value: attribute.StringValue(result)}) kv = append(t.kv, kv...) + opt := metric.WithAttributeSet(attribute.NewSet(kv...)) + t.RecordRaw(ctx, total, opt) +} + +// PrecomputeAttrs builds a reusable MeasurementOption from the given attribute +// key-values. The option must include all attributes (including "result"). +// Use with Stopwatch.Record to avoid per-call attribute allocation. +func PrecomputeAttrs(kv ...attribute.KeyValue) metric.MeasurementOption { + return metric.WithAttributeSet(attribute.NewSet(kv...)) +} +// RecordRaw records an operation using a precomputed attribute option, it does +// not include any previous attributes passed at Begin(). Zero-allocation +// alternative to Success/Failure for hot paths. +func (t Stopwatch) RecordRaw(ctx context.Context, total int64, precomputedAttrs metric.MeasurementOption) { amount := time.Since(t.start).Milliseconds() - t.histogram.Record(ctx, amount, metric.WithAttributes(kv...)) - t.sum.Add(ctx, total, metric.WithAttributes(kv...)) - t.count.Add(ctx, 1, metric.WithAttributes(kv...)) + t.histogram.Record(ctx, amount, precomputedAttrs) + t.sum.Add(ctx, total, precomputedAttrs) + t.count.Add(ctx, 1, precomputedAttrs) }