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
8 changes: 4 additions & 4 deletions packages/orchestrator/cmd/copy-build/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,13 +77,13 @@ func NewDestinationFromPath(prefix, file string) (*Destination, error) {
}, nil
}

func NewHeaderFromObject(ctx context.Context, bucketName string, headerPath string, objectType storage.ObjectType) (*header.Header, error) {
func NewHeaderFromObject(ctx context.Context, bucketName string, headerPath string) (*header.Header, error) {
b, err := storage.NewGCP(ctx, bucketName, nil)
if err != nil {
return nil, fmt.Errorf("failed to create GCS bucket storage provider: %w", err)
}

obj, err := b.OpenBlob(ctx, headerPath, objectType)
obj, err := b.OpenBlob(ctx, headerPath)
if err != nil {
return nil, fmt.Errorf("failed to open object: %w", err)
}
Expand Down Expand Up @@ -228,7 +228,7 @@ func main() {
if strings.HasPrefix(*from, "gs://") {
bucketName, _ := strings.CutPrefix(*from, "gs://")

h, err := NewHeaderFromObject(ctx, bucketName, buildMemfileHeaderPath, storage.MemfileHeaderObjectType)
h, err := NewHeaderFromObject(ctx, bucketName, buildMemfileHeaderPath)
if err != nil {
log.Fatalf("failed to create header from object: %s", err)
}
Expand All @@ -254,7 +254,7 @@ func main() {
var rootfsHeader *header.Header
if strings.HasPrefix(*from, "gs://") {
bucketName, _ := strings.CutPrefix(*from, "gs://")
h, err := NewHeaderFromObject(ctx, bucketName, buildRootfsHeaderPath, storage.RootFSHeaderObjectType)
h, err := NewHeaderFromObject(ctx, bucketName, buildRootfsHeaderPath)
if err != nil {
log.Fatalf("failed to create header from object: %s", err)
}
Expand Down
2 changes: 1 addition & 1 deletion packages/orchestrator/cmd/create-build/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -432,7 +432,7 @@ func printArtifactSizes(ctx context.Context, persistence storage.StorageProvider
printLocalFileSizes(basePath, buildID)
} else {
// For remote storage, get sizes from storage provider
if memfile, err := persistence.OpenSeekable(ctx, paths.Memfile(), storage.MemfileObjectType); err == nil {
if memfile, err := persistence.OpenSeekable(ctx, paths.Memfile()); err == nil {
if size, err := memfile.Size(ctx); err == nil {
fmt.Printf(" Memfile: %d MB\n", size>>20)
}
Expand Down
4 changes: 2 additions & 2 deletions packages/orchestrator/cmd/inspect-build/validate.go
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ func openChunker(ctx context.Context, storagePath, buildID, artifact string, h *
if ft.IsCompressed() {
dataPath += ft.CompressionType().Suffix()
}
obj, err := provider.OpenSeekable(ctx, dataPath, seekableType(artifact))
obj, err := provider.OpenSeekable(ctx, dataPath)
if err != nil {
return nil, nil, 0, nil, fmt.Errorf("open data: %w", err)
}
Expand Down Expand Up @@ -190,7 +190,7 @@ func openChunker(ctx context.Context, storagePath, buildID, artifact string, h *
return nil, nil, 0, nil, err
}

chunker, err := block.NewChunker(flags, size, int64(h.Metadata.BlockSize), filepath.Join(cacheDir, "cache"), m)
chunker, err := block.NewChunker(flags, size, int64(h.Metadata.BlockSize), filepath.Join(cacheDir, "cache"), m, seekableType(artifact))
if err != nil {
os.RemoveAll(cacheDir)

Expand Down
13 changes: 13 additions & 0 deletions packages/orchestrator/pkg/sandbox/block/fetch_session.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import (
"fmt"
"sync"
"sync/atomic"

"github.com/e2b-dev/infra/packages/shared/pkg/storage"
)

type fetchSession struct {
Expand All @@ -21,12 +23,23 @@ type fetchSession struct {
fetchErr error
done bool // true once terminated (success or error)

// source is the backend that served the fetch; set by runFetch.
source atomic.Int32

// bytesReady is the byte count (from chunkOff) up to which all blocks
// are fully written and marked cached. Atomic so registerAndWait can
// do a lock-free fast-path check: bytesReady only increases.
bytesReady atomic.Int64
}

func (s *fetchSession) setSource(src storage.Source) {
s.source.Store(int32(src))
}

func (s *fetchSession) Source() storage.Source {
return storage.Source(s.source.Load())
}

// contains reports whether the session covers the byte range [off, off+length).
func (s *fetchSession) contains(off, length int64) bool {
return s.chunkOff <= off && s.chunkOff+s.chunkLen >= off+length
Expand Down
44 changes: 7 additions & 37 deletions packages/orchestrator/pkg/sandbox/block/metrics/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,21 +8,10 @@ import (
"github.com/e2b-dev/infra/packages/shared/pkg/telemetry"
)

const (
orchestratorBlockSlices = "orchestrator.blocks.slices"
orchestratorBlockChunksFetch = "orchestrator.blocks.chunks.fetch"
orchestratorBlockChunksStore = "orchestrator.blocks.chunks.store"
)
const orchestratorChunkSlice = "orchestrator.chunk.slice"

type Metrics struct {
// SlicesMetric is used to measure page faulting performance.
SlicesTimerFactory telemetry.TimerFactory

// WriteChunksMetric is used to measure the time taken to download chunks from remote storage
RemoteReadsTimerFactory telemetry.TimerFactory

// WriteChunksMetric is used to measure performance of writing chunks to disk.
WriteChunksTimerFactory telemetry.TimerFactory
ChunkSliceTimerFactory telemetry.FloatTimerFactory
}

func NewMetrics(meterProvider metric.MeterProvider) (Metrics, error) {
Expand All @@ -31,31 +20,12 @@ func NewMetrics(meterProvider metric.MeterProvider) (Metrics, error) {
blocksMeter := meterProvider.Meter("github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block/metrics")

var err error
if m.SlicesTimerFactory, err = telemetry.NewTimerFactory(
blocksMeter, orchestratorBlockSlices,
"Time taken to retrieve memory slices",
"Total bytes requested",
"Total page faults",
); err != nil {
return m, fmt.Errorf("error creating slices timer factory: %w", err)
}

if m.RemoteReadsTimerFactory, err = telemetry.NewTimerFactory(
blocksMeter, orchestratorBlockChunksFetch,
"Time taken to fetch memory chunks from remote store",
"Total bytes fetched from remote store",
"Total remote fetches",
); err != nil {
return m, fmt.Errorf("error creating reads timer factory: %w", err)
}

if m.WriteChunksTimerFactory, err = telemetry.NewTimerFactory(
blocksMeter, orchestratorBlockChunksStore,
"Time taken to write memory chunks to disk",
"Total bytes written to disk",
"Total cache writes",
if m.ChunkSliceTimerFactory, err = telemetry.NewFloatTimerFactory(
blocksMeter, orchestratorChunkSlice,
"Time taken by Chunker to serve a Slice() (source=mmap when served from cache)",
"Bytes returned",
); err != nil {
return m, fmt.Errorf("failed to get stored chunks metric: %w", err)
return m, fmt.Errorf("error creating chunk slice timer factory: %w", err)
}

return m, nil
Expand Down
Loading
Loading