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
4 changes: 2 additions & 2 deletions packages/orchestrator/pkg/sandbox/block/cache_dirty_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -282,9 +282,9 @@ func TestSetIsCached_ConcurrentOverlapping(t *testing.T) {

const (
benchBlockSize int64 = 4096
benchNumBlocks int64 = 16384 // 64 MiB at 4K blocks — realistic memfile size
benchNumBlocks int64 = 16384 // 64 MiB at 4K blocks — realistic memfile size
benchCacheSize int64 = benchNumBlocks * benchBlockSize
benchChunkSize int64 = 4 * 1024 * 1024 // 4 MiB — MemoryChunkSize
benchChunkSize int64 = 4 * 1024 * 1024 // 4 MiB — MemoryChunkSize
benchChunkCount int64 = benchCacheSize / benchChunkSize
)

Expand Down
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) {
amount := time.Since(t.start).Milliseconds()
Comment thread
dobrac marked this conversation as resolved.
Comment thread
dobrac marked this conversation as resolved.
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