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
101 changes: 68 additions & 33 deletions packages/orchestrator/pkg/sandbox/block/chunk.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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)
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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
})
Expand All @@ -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"
)
60 changes: 60 additions & 0 deletions packages/orchestrator/pkg/sandbox/block/chunk_bench_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
}
24 changes: 7 additions & 17 deletions packages/orchestrator/pkg/sandbox/block/streaming_chunk.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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()
}

Expand Down
26 changes: 23 additions & 3 deletions packages/shared/pkg/telemetry/meters.go
Original file line number Diff line number Diff line change
Expand Up @@ -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...)
}
Expand All @@ -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) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

RecordRaw silently drops t.kv (the base attributes stored by Begin(kv ...)). All current call sites pass no args to Begin(), so this is fine today. But since RecordRaw is now exported, a future caller could do timer := factory.Begin(sandboxID, ...) and then call RecordRaw, silently losing those base dimensions from all three metric recordings without any compile-time or runtime warning.

Consider either: (a) merging t.kv into the precomputed option inside RecordRaw when t.kv is non-nil, or (b) keeping RecordRaw unexported and only exposing it through the typed Success/Failure methods.

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)
}
Loading