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
65 changes: 64 additions & 1 deletion packages/api/internal/orchestrator/create_instance.go
Original file line number Diff line number Diff line change
Expand Up @@ -318,15 +318,21 @@ func (o *Orchestrator) CreateSandbox(

allLabels, labelFilteringEnabled := o.generateRequiredNodeLabels(ctx, sandboxID, team, sbxData)

node, err = placement.PlaceSandbox(ctx, o.placementAlgorithm, clusterNodes, node, sbxRequest, builds.ToMachineInfo(sbxData.Build), labelFilteringEnabled, allLabels)
placed, err := placement.PlaceSandbox(ctx, o.placementAlgorithm, clusterNodes, node, sbxRequest, builds.ToMachineInfo(sbxData.Build), labelFilteringEnabled, allLabels)
if err != nil {
if isResume && placed.TimedOut {
o.maybeRemapResumeOriginNode(ctx, sandboxID, team, sbxData.NodeID, placed.WarmedNode)
}

return sandbox.Sandbox{}, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "Failed to place sandbox",
Err: fmt.Errorf("failed to place sandbox: %w", err),
}
}

node = placed.Node

// The sandbox was created successfully
attributes := []attribute.KeyValue{
attribute.Bool("is_resume", isResume),
Expand Down Expand Up @@ -409,6 +415,63 @@ func (o *Orchestrator) CreateSandbox(
return sbx, nil
}

// maybeRemapResumeOriginNode repoints the snapshot's origin_node_id to the
// fallback node a resume timed out on. Pinning the next resume attempt to avoid
// re-pulling the snapshot onto another node and spraying load across the cluster
//
// It only acts on placement timeouts (warmedNode is nil otherwise), and only
// when the warming node differs from the origin we already tried. The write runs
// on a detached context because the request context is already past its deadline
// (that is why placement timed out).
func (o *Orchestrator) maybeRemapResumeOriginNode(ctx context.Context, sandboxID string, team *teamtypes.Team, originNodeID *string, warmedNode *nodemanager.Node) {
if warmedNode == nil {
return
}

newNode := warmedNode
if originNodeID != nil && *originNodeID == newNode.ID {
return
Comment thread
jakubno marked this conversation as resolved.
}

// The request context is already past its deadline (that is why placement
// timed out), so detach it for everything below: the feature-flag read, the
// DB write, cache invalidation, the counter, and logging would all otherwise
// observe a cancelled context.
wctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer cancel()

if !o.featureFlagsClient.BoolFlag(wctx, featureflags.ResumeOriginNodeRemapFlag, featureflags.TeamContext(team.ID.String()), featureflags.SandboxContext(sandboxID)) {
return
}

if err := o.sqlcDB.UpdateSnapshotOriginNode(wctx, queries.UpdateSnapshotOriginNodeParams{
OriginNodeID: newNode.ID,
SandboxID: sandboxID,
}); err != nil {
logger.L().Warn(wctx, "failed to remap resume origin node",
zap.Error(err),
logger.WithSandboxID(sandboxID),
)

return
}

// Drop the cached snapshot so the next resume reads the new origin node.
o.snapshotCache.Invalidate(wctx, sandboxID)
o.resumeOriginNodeRemapCounter.Add(wctx, 1)

oldNodeID := ""
if originNodeID != nil {
oldNodeID = *originNodeID
}

logger.L().Info(wctx, "remapped resume origin node to the node a previous resume timed out on",
logger.WithSandboxID(sandboxID),
zap.String("old_origin_node_id", oldNodeID),
zap.String("new_origin_node_id", newNode.ID),
)
}

func (o *Orchestrator) generateRequiredNodeLabels(ctx context.Context, sandboxID string, team *teamtypes.Team, sbxData SandboxMetadata) ([]string, bool) {
labelFilteringEnabled := o.featureFlagsClient.BoolFlag(ctx, featureflags.SandboxLabelBasedSchedulingFlag, featureflags.TeamContext(team.ID.String()), featureflags.SandboxContext(sandboxID))
if !labelFilteringEnabled {
Expand Down
4 changes: 4 additions & 0 deletions packages/api/internal/orchestrator/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,10 @@ func (o *Orchestrator) setupMetrics(meterProvider metric.MeterProvider) error {
return fmt.Errorf("failed to create sandboxes counter: %w", err)
}

if o.resumeOriginNodeRemapCounter, err = telemetry.GetCounter(meter, telemetry.ApiOrchestratorResumeOriginNodeRemap); err != nil {
return fmt.Errorf("failed to create resume origin node remap counter: %w", err)
}

// Observable gauge that reads sandbox counts from Redis on each collection interval.
// This replaces the old UpDownCounter which drifted across multiple API instances.
sandboxCountGauge, err := telemetry.GetGaugeInt(meter, telemetry.SandboxCountGaugeName)
Expand Down
1 change: 1 addition & 0 deletions packages/api/internal/orchestrator/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ type Orchestrator struct {
metricsRegistration metric.Registration
sandboxCountGaugeRegistration metric.Registration
createdSandboxesCounter metric.Int64Counter
resumeOriginNodeRemapCounter metric.Int64Counter
teamMetricsObserver *metrics.TeamObserver
accessTokenGenerator *sandbox.AccessTokenGenerator
createdCounter metric.Int64Counter
Expand Down
49 changes: 43 additions & 6 deletions packages/api/internal/orchestrator/placement/placement.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,17 @@ var tracer = otel.Tracer("github.com/e2b-dev/infra/packages/api/internal/orchest

var errSandboxCreateFailed = errors.New("failed to create a new sandbox, if the problem persists, contact us")

// PlacementResult carries the outcome of a placement attempt alongside the error.
type PlacementResult struct {
// Node is the node the sandbox was placed on, or nil on failure.
Node *nodemanager.Node
// WarmedNode is the first node that attempted the create before a timeout,
// whose cache is warmest; nil unless TimedOut.
WarmedNode *nodemanager.Node
// TimedOut reports whether placement failed due to context cancellation/deadline.
TimedOut bool
}

// Algorithm defines the interface for sandbox placement strategies.
// Implementations should choose an optimal node based on available resources
// and current load distribution.
Expand All @@ -38,7 +49,7 @@ func PlaceSandbox(
buildMachineInfo machineinfo.MachineInfo,
labelFilteringEnabled bool,
requiredLabels []string,
) (*nodemanager.Node, error) {
) (PlacementResult, error) {
ctx, span := tracer.Start(ctx, "place-sandbox")
defer span.End()

Expand All @@ -50,11 +61,31 @@ func PlaceSandbox(
node = preferredNode
}

// First node that attempted the create (not a fast ResourceExhausted refusal).
var firstTriedNode *nodemanager.Node

// failed reports the warming node only when the failure was caused by the
// request context being cancelled or timing out (ctx.Err() != nil). Hard
// failures (where the context is still live) carry no node, so callers never
// pin a retry to a node that genuinely refused the sandbox.
//
// TODO [EN-1099]: We key off ctx.Err() rather than the gRPC status code because
// the orchestrator currently collapses a timed-out resume into codes.Internal
// (it folds the deadline cause into the message, not the code),
// so the code alone cannot tell a timeout apart from a hard failure.
failed := func(err error) (PlacementResult, error) {
if ctx.Err() == nil {
return PlacementResult{}, err
}

return PlacementResult{WarmedNode: firstTriedNode, TimedOut: true}, err
}

attempt := 0
for attempt < maxRetries {
select {
case <-ctx.Done():
return nil, fmt.Errorf("request timed out during %d. attempt", attempt+1)
return failed(fmt.Errorf("request timed out during %d. attempt", attempt+1))
default:
// Continue
}
Expand All @@ -63,12 +94,12 @@ func PlaceSandbox(
telemetry.ReportEvent(ctx, "Placing sandbox on the preferred node", telemetry.WithNodeID(node.ID))
} else {
if len(nodesExcluded) >= len(clusterNodes) {
return nil, errors.New("no nodes available")
return failed(errors.New("no nodes available"))
}

node, err = algorithm.chooseNode(ctx, clusterNodes, nodesExcluded, nodemanager.SandboxResources{CPUs: sbxRequest.GetSandbox().GetVcpu(), MiBMemory: sbxRequest.GetSandbox().GetRamMb()}, buildMachineInfo, labelFilteringEnabled, requiredLabels)
if err != nil {
return nil, err
return failed(err)
}

telemetry.ReportEvent(ctx, "Placing sandbox on the node", telemetry.WithNodeID(node.ID))
Expand Down Expand Up @@ -97,7 +128,7 @@ func PlaceSandbox(
MiBMemory: sbxRequest.GetSandbox().GetRamMb(),
})

return node, nil
return PlacementResult{Node: node}, nil
}

failedNode := node
Expand All @@ -109,6 +140,12 @@ func PlaceSandbox(
statusCode = st.Code()
}

// Remember the first node that got far enough to actually attempt the
// sandbox (i.e. did not refuse with ResourceExhausted).
if statusCode != codes.ResourceExhausted && firstTriedNode == nil {
firstTriedNode = failedNode
Comment thread
jakubno marked this conversation as resolved.
}
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
jakubno marked this conversation as resolved.

switch statusCode {
case codes.ResourceExhausted:
failedNode.PlacementMetrics.Skip(sbxRequest.GetSandbox().GetSandboxId())
Expand All @@ -121,5 +158,5 @@ func PlaceSandbox(
}
}

return nil, errSandboxCreateFailed
return failed(errSandboxCreateFailed)
}
Original file line number Diff line number Diff line change
Expand Up @@ -419,7 +419,7 @@ func runBenchmark(b *testing.B, algorithm Algorithm, config BenchmarkConfig, nod
wg.Go(func(sbx *LiveSandbox) func() {
return func() {
placementStart := time.Now()
node, err := PlaceSandbox(ctx, algorithm, nodes, nil, &orchestratorgrpc.SandboxCreateRequest{Sandbox: &orchestratorgrpc.SandboxConfig{
result, err := PlaceSandbox(ctx, algorithm, nodes, nil, &orchestratorgrpc.SandboxCreateRequest{Sandbox: &orchestratorgrpc.SandboxConfig{
SandboxId: sbx.ID,
Vcpu: sbx.RequestedCPU,
RamMb: sbx.RequestedMemory,
Expand All @@ -441,10 +441,10 @@ func runBenchmark(b *testing.B, algorithm Algorithm, config BenchmarkConfig, nod
}

success := false
if err == nil && node != nil {
if err == nil && result.Node != nil {
// Find the simulated node and place the sandbox
if simNode, exists := nodeMap[node.ID]; exists {
sbx.NodeID = node.ID
if simNode, exists := nodeMap[result.Node.ID]; exists {
sbx.NodeID = result.Node.ID
if simNode.PlaceSandbox(sbx) {
activeSandboxes.Store(sbx.ID, sbx)
metrics.SuccessfulPlacements++
Expand Down Expand Up @@ -717,16 +717,16 @@ func BenchmarkPlacementDistribution(b *testing.B) {
wg.Add(1)
go func(s *LiveSandbox) {
// Execute placement algorithm
node, err := PlaceSandbox(ctx, alg.algo, nodes, nil, &orchestratorgrpc.SandboxCreateRequest{
result, err := PlaceSandbox(ctx, alg.algo, nodes, nil, &orchestratorgrpc.SandboxCreateRequest{
Sandbox: &orchestratorgrpc.SandboxConfig{
SandboxId: s.ID,
Vcpu: s.RequestedCPU,
RamMb: s.RequestedMemory,
},
}, machineinfo.MachineInfo{}, false, nil)

if err == nil && node != nil {
if simNode, ok := nodeMap[node.ID]; ok {
if err == nil && result.Node != nil {
if simNode, ok := nodeMap[result.Node.ID]; ok {
// Placement successful, but Metrics won't update immediately (LaggyNode feature)
simNode.PlaceSandbox(s)
}
Expand Down
26 changes: 13 additions & 13 deletions packages/api/internal/orchestrator/placement/placement_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,8 @@ func TestPlaceSandbox_SuccessfulPlacement(t *testing.T) {
resultNode, err := PlaceSandbox(ctx, algorithm, nodes, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.NoError(t, err)
assert.NotNil(t, resultNode)
assert.Equal(t, node2, resultNode)
assert.NotNil(t, resultNode.Node)
assert.Equal(t, node2, resultNode.Node)
algorithm.AssertExpectations(t)
}

Expand Down Expand Up @@ -90,15 +90,15 @@ func TestPlaceSandbox_WithPreferredNode(t *testing.T) {

resultNode, err := PlaceSandbox(ctx, algorithm, nodes, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)
require.NoError(t, err)
assert.NotNil(t, resultNode)
assert.Equal(t, node1, resultNode)
assert.NotNil(t, resultNode.Node)
assert.Equal(t, node1, resultNode.Node)
algorithm.AssertExpectations(t)

// Test with preferred node - should use the preferred node directly without calling algorithm
resultNode, err = PlaceSandbox(ctx, algorithm, nodes, node2, sbxRequest, machineinfo.MachineInfo{}, false, nil)
require.NoError(t, err)
assert.NotNil(t, resultNode)
assert.Equal(t, node2, resultNode)
assert.NotNil(t, resultNode.Node)
assert.Equal(t, node2, resultNode.Node)
// Algorithm should not be called when preferred node is provided
algorithm.AssertNotCalled(t, "chooseNode")
}
Expand Down Expand Up @@ -129,7 +129,7 @@ func TestPlaceSandbox_ContextTimeout(t *testing.T) {
}, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.Error(t, err)
assert.Nil(t, resultNode)
assert.Nil(t, resultNode.Node)
// The error could be either "timeout" from the algorithm or "request timed out" from ctx.Done()
assert.True(t, err.Error() == "timeout" || strings.Contains(err.Error(), "request timed out"))
}
Expand All @@ -150,7 +150,7 @@ func TestPlaceSandbox_NoNodes(t *testing.T) {
resultNode, err := PlaceSandbox(ctx, algorithm, []*nodemanager.Node{}, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.Error(t, err)
assert.Nil(t, resultNode)
assert.Nil(t, resultNode.Node)
assert.Contains(t, err.Error(), "no nodes available")
}

Expand All @@ -175,7 +175,7 @@ func TestPlaceSandbox_AllNodesExcluded(t *testing.T) {
}, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.Error(t, err)
assert.Nil(t, resultNode)
assert.Nil(t, resultNode.Node)
assert.Contains(t, err.Error(), "no nodes available")
algorithm.AssertExpectations(t)
}
Expand Down Expand Up @@ -208,8 +208,8 @@ func TestPlaceSandbox_ResourceExhausted(t *testing.T) {
resultNode, err := PlaceSandbox(ctx, algorithm, nodes, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.NoError(t, err)
assert.NotNil(t, resultNode)
assert.Equal(t, node2, resultNode, "should succeed on node2 after node1 was exhausted")
assert.NotNil(t, resultNode.Node)
assert.Equal(t, node2, resultNode.Node, "should succeed on node2 after node1 was exhausted")
algorithm.AssertExpectations(t)

// Verify node1 was NOT excluded (ResourceExhausted nodes should be retried)
Expand Down Expand Up @@ -249,9 +249,9 @@ func TestPlaceSandbox_TriggersOptimisticUpdate(t *testing.T) {
resultNode, err := PlaceSandbox(ctx, algorithm, nodes, nil, sbxRequest, machineinfo.MachineInfo{}, false, nil)

require.NoError(t, err)
assert.NotNil(t, resultNode)
assert.NotNil(t, resultNode.Node)

// Verify: After successful placement, the node's CpuAllocated should be increased by 2 from the base
updatedCpuAllocated := resultNode.Metrics().CpuAllocated
updatedCpuAllocated := resultNode.Node.Metrics().CpuAllocated
assert.Equal(t, initialCpuAllocated+2, updatedCpuAllocated, "Node metrics should be optimistically updated after placement")
}
Loading
Loading