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
75 changes: 56 additions & 19 deletions packages/api/internal/orchestrator/evictor/evict.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,40 +10,54 @@ import (
"github.com/google/uuid"
"go.opentelemetry.io/otel/metric"
"go.uber.org/zap"
"golang.org/x/sync/errgroup"

"github.com/e2b-dev/infra/packages/api/internal/pause"
"github.com/e2b-dev/infra/packages/api/internal/sandbox"
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
"github.com/e2b-dev/infra/packages/shared/pkg/telemetry"
"github.com/e2b-dev/infra/packages/shared/pkg/utils"
)

const (
pollInterval = 50 * time.Millisecond

// maxConcurrentEvictions caps the number of evictions that can run in
// parallel. Excess items remain expired in the store and are picked up by
// the next tick.
maxConcurrentEvictions = 256
pollInterval = 50 * time.Millisecond
concurrencyRefreshInterval = 30 * time.Second
)

type Evictor struct {
store *sandbox.Store
removeSandbox func(ctx context.Context, teamID uuid.UUID, sandboxID string, opts sandbox.RemoveOpts) error
featureFlags *featureflags.Client

concurrencyLimiter *utils.AdjustableSemaphore

// activeEvictions tracks concurrent eviction attempts for the same sandbox
// so that overlapping ticks don't kick off multiple removeSandbox calls.
activeEvictions sync.Map
}

func New(
ctx context.Context,
store *sandbox.Store,
removeSandbox func(ctx context.Context, teamID uuid.UUID, sandboxID string, opts sandbox.RemoveOpts) error,
featureFlags *featureflags.Client,
meter metric.Meter,
) (*Evictor, error) {
initialLimit := featureFlags.IntFlag(ctx, featureflags.MaxConcurrentEvictions)
if initialLimit <= 0 {
initialLimit = featureflags.MaxConcurrentEvictions.Fallback()
}

concurrencyLimiter, err := utils.NewAdjustableSemaphore(int64(initialLimit))
if err != nil {
return nil, fmt.Errorf("failed to create eviction concurrency semaphore: %w", err)
}

e := &Evictor{
store: store,
removeSandbox: removeSandbox,
store: store,
removeSandbox: removeSandbox,
featureFlags: featureFlags,
concurrencyLimiter: concurrencyLimiter,
}

if _, err := telemetry.GetObservableUpDownCounter(meter, telemetry.EvictionsRunningCounterName,
Expand All @@ -66,18 +80,22 @@ func New(
}

func (e *Evictor) Start(ctx context.Context) {
g := errgroup.Group{}
g.SetLimit(maxConcurrentEvictions)
var wg sync.WaitGroup
ticker := time.NewTicker(pollInterval)
defer ticker.Stop()

refreshTicker := time.NewTicker(concurrencyRefreshInterval)
defer refreshTicker.Stop()

for {
select {
case <-ctx.Done():
// Wait for in-flight evictions to finish for graceful shutdown.
g.Wait()
wg.Wait()

return
case <-refreshTicker.C:
e.refreshConcurrencyLimit(ctx)
case <-ticker.C:
sbxs, err := e.store.ExpiredItems(ctx)
if err != nil {
Expand All @@ -92,25 +110,44 @@ func (e *Evictor) Start(ctx context.Context) {
continue
}

if ok := g.TryGo(func() error {
defer e.activeEvictions.Delete(item.SandboxID)

e.evictSandbox(ctx, item)

return nil
}); !ok {
// Non-blocking acquire: if we're at capacity, skip and let the
// next tick retry. Mirrors the previous errgroup.TryGo behavior.
if !e.concurrencyLimiter.TryAcquire(1) {
e.activeEvictions.Delete(item.SandboxID)

logger.L().Debug(ctx, "Max concurrent evictions reached, skipping eviction this tick",
logger.WithSandboxID(item.SandboxID),
logger.WithTeamID(item.TeamID.String()),
)

continue
}

wg.Add(1)
go func(item sandbox.Sandbox) {
defer wg.Done()
defer e.concurrencyLimiter.Release(1)
defer e.activeEvictions.Delete(item.SandboxID)

e.evictSandbox(ctx, item)
}(item)
Comment thread
jakubno marked this conversation as resolved.
}
}
}
}

func (e *Evictor) refreshConcurrencyLimit(ctx context.Context) {
limit := e.featureFlags.IntFlag(ctx, featureflags.MaxConcurrentEvictions)
if limit <= 0 {
return
}

if err := e.concurrencyLimiter.SetLimit(int64(limit)); err != nil {
logger.L().Error(ctx, "failed to adjust eviction concurrency semaphore",
zap.Int("limit", limit), zap.Error(err))
}
}

func (e *Evictor) evictSandbox(ctx context.Context, sbx sandbox.Sandbox) {
action := sandbox.StateActionKill
if sbx.AutoPause {
Expand Down
2 changes: 1 addition & 1 deletion packages/api/internal/orchestrator/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ func New(
)

// Evict old sandboxes
sandboxEvictor, err := evictor.New(o.sandboxStore, o.RemoveSandbox, meter)
sandboxEvictor, err := evictor.New(ctx, o.sandboxStore, o.RemoveSandbox, o.featureFlagsClient, meter)
if err != nil {
return nil, fmt.Errorf("failed to create sandbox evictor: %w", err)
}
Expand Down
6 changes: 6 additions & 0 deletions packages/shared/pkg/featureflags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -211,6 +211,12 @@ var (
// Must be > 0.
MaxStartingInstancesPerNode = NewIntFlag("max-starting-instances-per-node", 3)

// MaxConcurrentEvictions caps the number of sandbox evictions that can run
// in parallel per API instance. Excess items remain expired in the store
// and are picked up by the next eviction tick. Must be > 0; non-positive
// values are ignored at refresh time.
MaxConcurrentEvictions = NewIntFlag("max-concurrent-evictions", 256)

// MaxConcurrentSnapshotUpserts limits concurrent UpsertSnapshot calls (pause + snapshot template paths).
// 0 or negative disables throttling (unlimited concurrency).
MaxConcurrentSnapshotUpserts = NewIntFlag("max-concurrent-snapshot-upserts", 0)
Expand Down
Loading