diff --git a/packages/api/internal/orchestrator/create_instance.go b/packages/api/internal/orchestrator/create_instance.go index 41677dfe49..da26c66da9 100644 --- a/packages/api/internal/orchestrator/create_instance.go +++ b/packages/api/internal/orchestrator/create_instance.go @@ -318,8 +318,12 @@ 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", @@ -327,6 +331,8 @@ func (o *Orchestrator) CreateSandbox( } } + node = placed.Node + // The sandbox was created successfully attributes := []attribute.KeyValue{ attribute.Bool("is_resume", isResume), @@ -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 + } + + // 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 { diff --git a/packages/api/internal/orchestrator/metrics.go b/packages/api/internal/orchestrator/metrics.go index fb15f75daa..7e70dc94d9 100644 --- a/packages/api/internal/orchestrator/metrics.go +++ b/packages/api/internal/orchestrator/metrics.go @@ -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) diff --git a/packages/api/internal/orchestrator/orchestrator.go b/packages/api/internal/orchestrator/orchestrator.go index 472fb34128..7ceb5b33a1 100644 --- a/packages/api/internal/orchestrator/orchestrator.go +++ b/packages/api/internal/orchestrator/orchestrator.go @@ -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 diff --git a/packages/api/internal/orchestrator/placement/placement.go b/packages/api/internal/orchestrator/placement/placement.go index b2bbb8a0c1..4ad0ef292a 100644 --- a/packages/api/internal/orchestrator/placement/placement.go +++ b/packages/api/internal/orchestrator/placement/placement.go @@ -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. @@ -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() @@ -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 } @@ -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)) @@ -97,7 +128,7 @@ func PlaceSandbox( MiBMemory: sbxRequest.GetSandbox().GetRamMb(), }) - return node, nil + return PlacementResult{Node: node}, nil } failedNode := node @@ -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 + } + switch statusCode { case codes.ResourceExhausted: failedNode.PlacementMetrics.Skip(sbxRequest.GetSandbox().GetSandboxId()) @@ -121,5 +158,5 @@ func PlaceSandbox( } } - return nil, errSandboxCreateFailed + return failed(errSandboxCreateFailed) } diff --git a/packages/api/internal/orchestrator/placement/placement_benchmark_test.go b/packages/api/internal/orchestrator/placement/placement_benchmark_test.go index 3f57761a98..1cf547c07b 100644 --- a/packages/api/internal/orchestrator/placement/placement_benchmark_test.go +++ b/packages/api/internal/orchestrator/placement/placement_benchmark_test.go @@ -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, @@ -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++ @@ -717,7 +717,7 @@ 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, @@ -725,8 +725,8 @@ func BenchmarkPlacementDistribution(b *testing.B) { }, }, 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) } diff --git a/packages/api/internal/orchestrator/placement/placement_test.go b/packages/api/internal/orchestrator/placement/placement_test.go index 3d27d43a8d..cbc5cd4f04 100644 --- a/packages/api/internal/orchestrator/placement/placement_test.go +++ b/packages/api/internal/orchestrator/placement/placement_test.go @@ -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) } @@ -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") } @@ -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")) } @@ -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") } @@ -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) } @@ -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) @@ -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") } diff --git a/packages/api/internal/orchestrator/placement/placement_timeout_test.go b/packages/api/internal/orchestrator/placement/placement_timeout_test.go new file mode 100644 index 0000000000..15054f175a --- /dev/null +++ b/packages/api/internal/orchestrator/placement/placement_timeout_test.go @@ -0,0 +1,233 @@ +package placement + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + + "github.com/e2b-dev/infra/packages/api/internal/api" + "github.com/e2b-dev/infra/packages/api/internal/orchestrator/nodemanager" + "github.com/e2b-dev/infra/packages/shared/pkg/grpc/orchestrator" + "github.com/e2b-dev/infra/packages/shared/pkg/machineinfo" +) + +// stubAlgorithm is a placement Algorithm whose chooseNode behavior is injected. +type stubAlgorithm struct { + choose func(nodesExcluded map[string]struct{}) (*nodemanager.Node, error) +} + +func (s stubAlgorithm) chooseNode( + _ context.Context, + _ []*nodemanager.Node, + nodesExcluded map[string]struct{}, + _ nodemanager.SandboxResources, + _ machineinfo.MachineInfo, + _ bool, + _ []string, +) (*nodemanager.Node, error) { + return s.choose(nodesExcluded) +} + +func testSbxRequest(id string) *orchestrator.SandboxCreateRequest { + return &orchestrator.SandboxCreateRequest{ + Sandbox: &orchestrator.SandboxConfig{ + SandboxId: id, + Vcpu: 2, + RamMb: 512, + }, + } +} + +// failIfCalled is an algorithm that fails the test if chooseNode is invoked. +func failIfCalled(t *testing.T) stubAlgorithm { + t.Helper() + + return stubAlgorithm{ + choose: func(map[string]struct{}) (*nodemanager.Node, error) { + t.Fatal("chooseNode should not be called") + + return nil, nil + }, + } +} + +// erroringClient returns err from SandboxCreate, optionally cancelling the +// request context first to simulate the deadline being hit mid-create. +func erroringClient(cancel context.CancelFunc, err error) *nodemanager.MockSandboxClientCustom { + return &nodemanager.MockSandboxClientCustom{ + CreateFunc: func() error { + if cancel != nil { + cancel() + } + + return err + }, + } +} + +// TestPlaceSandbox_TimeoutPinsFirstTriedNode is the core prod scenario: a node +// fails (here with codes.Internal, the code the orchestrator returns for a +// timed-out resume) and the request context is cancelled. PlaceSandbox must +// surface that node so a retry can be pinned to it - detection must NOT depend +// on the gRPC code being DeadlineExceeded/Canceled. +func TestPlaceSandbox_TimeoutPinsFirstTriedNode(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(t.Context()) + + node := nodemanager.NewTestNode("node-test", api.NodeStatusReady, 0, 8) + node.SetSandboxClient(erroringClient(cancel, status.Error(codes.Internal, "failed to create sandbox: request timed out"))) + + result, err := PlaceSandbox( + ctx, + failIfCalled(t), + []*nodemanager.Node{node}, + node, + testSbxRequest("sbx-1"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.Error(t, err) + assert.True(t, result.TimedOut, "expected the failure to be reported as a timeout") + require.NotNil(t, result.WarmedNode) + assert.Equal(t, node.ID, result.WarmedNode.ID) +} + +// TestPlaceSandbox_PinsFirstTriedNodeNotLater verifies that across multiple +// attempts the FIRST node tried is the one surfaced, not a later one. +func TestPlaceSandbox_PinsFirstTriedNodeNotLater(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(t.Context()) + + first := nodemanager.NewTestNode("node-first", api.NodeStatusReady, 0, 8) + // First node fails without cancelling - a long pull that the orchestrator + // gave up on while the request budget was still alive. + first.SetSandboxClient(erroringClient(nil, status.Error(codes.Internal, "boom"))) + + second := nodemanager.NewTestNode("node-second", api.NodeStatusReady, 0, 8) + // Second node fails and the request budget runs out here. + second.SetSandboxClient(erroringClient(cancel, status.Error(codes.Internal, "boom"))) + + algo := stubAlgorithm{choose: func(map[string]struct{}) (*nodemanager.Node, error) { + return second, nil + }} + + result, err := PlaceSandbox( + ctx, + algo, + []*nodemanager.Node{first, second}, + first, + testSbxRequest("sbx-2"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.Error(t, err) + assert.True(t, result.TimedOut) + require.NotNil(t, result.WarmedNode) + assert.Equal(t, first.ID, result.WarmedNode.ID, "must pin the first node tried, not a later one") +} + +// TestPlaceSandbox_HardFailureNotWrapped verifies a failure where the context is +// still live (a genuine error, not a timeout) is returned unchanged, so a retry +// is never pinned to a node that actually refused the sandbox. +func TestPlaceSandbox_HardFailureNotWrapped(t *testing.T) { + t.Parallel() + + node := nodemanager.NewTestNode("node-internal", api.NodeStatusReady, 0, 8) + node.SetSandboxClient(erroringClient(nil, status.Error(codes.Internal, "boom"))) + + result, err := PlaceSandbox( + t.Context(), + failIfCalled(t), + []*nodemanager.Node{node}, + node, + testSbxRequest("sbx-3"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.Error(t, err) + assert.False(t, result.TimedOut, "a live-context failure must not be reported as a timeout") + assert.Nil(t, result.WarmedNode, "a live-context failure must not surface a warming node") +} + +// TestPlaceSandbox_ResourceExhaustedNotPinned verifies that a node which refused +// fast with ResourceExhausted is not pinned even if the request then times out. +func TestPlaceSandbox_ResourceExhaustedNotPinned(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(t.Context()) + + node := nodemanager.NewTestNode("node-exhausted", api.NodeStatusReady, 0, 8) + node.SetSandboxClient(erroringClient(cancel, status.Error(codes.ResourceExhausted, "no capacity"))) + + result, err := PlaceSandbox( + ctx, + failIfCalled(t), + []*nodemanager.Node{node}, + node, + testSbxRequest("sbx-4"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.Error(t, err) + assert.Nil(t, result.WarmedNode, "a node that refused fast must not be pinned") +} + +// TestPlaceSandbox_TimeoutBeforeAnyAttemptNotWrapped verifies that when the +// deadline fires before any node was tried there is nothing to pin to. +func TestPlaceSandbox_TimeoutBeforeAnyAttemptNotWrapped(t *testing.T) { + t.Parallel() + + node := nodemanager.NewTestNode("node-ctx", api.NodeStatusReady, 0, 8) + + ctx, cancel := context.WithCancel(t.Context()) + cancel() // already past deadline before the first attempt + + result, err := PlaceSandbox( + ctx, + failIfCalled(t), + []*nodemanager.Node{node}, + node, + testSbxRequest("sbx-5"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.Error(t, err) + assert.Nil(t, result.WarmedNode, "no node was tried, so there is nothing to pin") +} + +func TestPlaceSandbox_SuccessReturnsNode(t *testing.T) { + t.Parallel() + + node := nodemanager.NewTestNode("node-ok", api.NodeStatusReady, 0, 8) + + result, err := PlaceSandbox( + t.Context(), + failIfCalled(t), + []*nodemanager.Node{node}, + node, + testSbxRequest("sbx-6"), + machineinfo.MachineInfo{}, + false, + nil, + ) + + require.NoError(t, err) + require.NotNil(t, result.Node) + assert.Equal(t, node.ID, result.Node.ID) +} diff --git a/packages/db/queries/snapshots/update_snapshot_origin_node.sql b/packages/db/queries/snapshots/update_snapshot_origin_node.sql new file mode 100644 index 0000000000..290a412a0a --- /dev/null +++ b/packages/db/queries/snapshots/update_snapshot_origin_node.sql @@ -0,0 +1,6 @@ +-- name: UpdateSnapshotOriginNode :exec +-- Repoints a snapshot's origin node, used on resume to pin a retry to the node +-- whose local cache is warming after a previous resume attempt timed out. +UPDATE "public"."snapshots" +SET origin_node_id = @origin_node_id +WHERE sandbox_id = @sandbox_id; diff --git a/packages/db/queries/update_snapshot_origin_node.sql.go b/packages/db/queries/update_snapshot_origin_node.sql.go new file mode 100644 index 0000000000..7ecad677cc --- /dev/null +++ b/packages/db/queries/update_snapshot_origin_node.sql.go @@ -0,0 +1,28 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.29.0 +// source: update_snapshot_origin_node.sql + +package queries + +import ( + "context" +) + +const updateSnapshotOriginNode = `-- name: UpdateSnapshotOriginNode :exec +UPDATE "public"."snapshots" +SET origin_node_id = $1 +WHERE sandbox_id = $2 +` + +type UpdateSnapshotOriginNodeParams struct { + OriginNodeID string + SandboxID string +} + +// Repoints a snapshot's origin node, used on resume to pin a retry to the node +// whose local cache is warming after a previous resume attempt timed out. +func (q *Queries) UpdateSnapshotOriginNode(ctx context.Context, arg UpdateSnapshotOriginNodeParams) error { + _, err := q.db.Exec(ctx, updateSnapshotOriginNode, arg.OriginNodeID, arg.SandboxID) + return err +} diff --git a/packages/shared/pkg/featureflags/flags.go b/packages/shared/pkg/featureflags/flags.go index 11a3f9c8ff..5aa38004cd 100644 --- a/packages/shared/pkg/featureflags/flags.go +++ b/packages/shared/pkg/featureflags/flags.go @@ -213,6 +213,12 @@ var ( // HeaderV5WriteFlag makes Pause emit V5 headers. When enabled it also // supersedes V4HeaderForUncompressedFlag for uncompressed uploads. HeaderV5WriteFlag = NewBoolFlag("header-v5-write", false) + + // ResumeOriginNodeRemapFlag enables repointing a snapshot's origin_node_id to + // the fallback node a resume timed out on. The node's local cache is warming + // from the in-progress snapshot pull, so pinning the retry to it avoids + // re-pulling the snapshot onto yet another node. + ResumeOriginNodeRemapFlag = NewBoolFlag("resume-origin-node-remap", false) ) type IntFlag struct { diff --git a/packages/shared/pkg/telemetry/meters.go b/packages/shared/pkg/telemetry/meters.go index 2d77c8593e..24eb5850e0 100644 --- a/packages/shared/pkg/telemetry/meters.go +++ b/packages/shared/pkg/telemetry/meters.go @@ -20,8 +20,9 @@ type ( ) const ( - ApiOrchestratorCreatedSandboxes CounterType = "api.orchestrator.created_sandboxes" - SandboxCreateMeterName CounterType = "api.env.instance.started" + ApiOrchestratorCreatedSandboxes CounterType = "api.orchestrator.created_sandboxes" + ApiOrchestratorResumeOriginNodeRemap CounterType = "api.orchestrator.resume_origin_node_remapped" + SandboxCreateMeterName CounterType = "api.env.instance.started" TeamSandboxCreated CounterType = "e2b.team.sandbox.created" @@ -197,6 +198,7 @@ const ( var counterDesc = map[CounterType]string{ SandboxCreateMeterName: "Number of currently waiting requests to create a new sandbox", ApiOrchestratorCreatedSandboxes: "Number of successfully created sandboxes", + ApiOrchestratorResumeOriginNodeRemap: "Number of resume snapshots repointed to the fallback node a previous resume timed out on", BuildResultCounterName: "Number of template build results", BuildCacheResultCounterName: "Number of build cache results", TeamSandboxCreated: "Counter of started sandboxes for the team in the interval", @@ -226,6 +228,7 @@ var counterDesc = map[CounterType]string{ var counterUnits = map[CounterType]string{ SandboxCreateMeterName: "{sandbox}", ApiOrchestratorCreatedSandboxes: "{sandbox}", + ApiOrchestratorResumeOriginNodeRemap: "{snapshot}", BuildResultCounterName: "{build}", BuildCacheResultCounterName: "{layer}", TeamSandboxCreated: "{sandbox}",