diff --git a/packages/orchestrator/pkg/sandbox/map.go b/packages/orchestrator/pkg/sandbox/map.go index 0f2bed8f84..867a1f3ecb 100644 --- a/packages/orchestrator/pkg/sandbox/map.go +++ b/packages/orchestrator/pkg/sandbox/map.go @@ -179,6 +179,10 @@ var ErrSandboxAlreadyRunning = errors.New("sandbox is already running on this no var ErrSandboxOperationInProgress = errors.New("sandbox operation already in progress") +// ErrNodeAtCapacity reports that live sandboxes and held reservations already +// fill the limit passed to ReserveWithin. +var ErrNodeAtCapacity = errors.New("node is at its sandbox limit") + // Reservation prevents concurrent starts; a failed start retains it through its own cleanup. type Reservation struct { m *Map @@ -188,6 +192,24 @@ type Reservation struct { // Reserve takes the sandbox ID for a create that is about to start a VM. // Refused with ErrSandboxAlreadyRunning while the ID is live or reserved. func (m *Map) Reserve(sandboxID string) (*Reservation, error) { + return m.reserve(sandboxID, 0, false) +} + +// ReserveWithin is Reserve with admission control. The ID is taken only while +// the sandboxes this node holds stay below limit. A sandbox is held while its +// ID is live or reserved, so starts that have not reached MarkRunning, +// failed starts still in cleanup, and checkpoint holds all count. A reservation +// whose sandbox is already live counts once. +// +// Counting and inserting happen under the same lock, so concurrent creates +// cannot all see the same free slot. A non-positive limit refuses every create. +// Duplicate IDs are still refused with ErrSandboxAlreadyRunning, before the +// limit is checked, so the caller does not retry them on another node. +func (m *Map) ReserveWithin(sandboxID string, limit int64) (*Reservation, error) { + return m.reserve(sandboxID, limit, true) +} + +func (m *Map) reserve(sandboxID string, limit int64, bounded bool) (*Reservation, error) { m.registryMu.Lock() defer m.registryMu.Unlock() @@ -197,12 +219,30 @@ func (m *Map) Reserve(sandboxID string) (*Reservation, error) { if _, ok := m.reservations[sandboxID]; ok { return nil, fmt.Errorf("%w: another create is in flight", ErrSandboxAlreadyRunning) } + if bounded { + if held := m.heldLocked(); held >= limit { + return nil, fmt.Errorf("%w: %d held, limit %d", ErrNodeAtCapacity, held, limit) + } + } r := &Reservation{m: m, sandboxID: sandboxID} m.reservations[sandboxID] = r return r, nil } +// heldLocked counts live sandboxes plus reservations whose ID is not live. +// The caller must hold registryMu; every live insert and removal takes it too. +func (m *Map) heldLocked() int64 { + held := int64(m.live.Count()) + for id := range m.reservations { + if _, ok := m.live.Get(id); !ok { + held++ + } + } + + return held +} + // Release ends operation ownership; the live entry and physical lifecycles remain independent. func (r *Reservation) Release() { r.m.registryMu.Lock() diff --git a/packages/orchestrator/pkg/sandbox/map_test.go b/packages/orchestrator/pkg/sandbox/map_test.go index 720879681d..8948aeadbe 100644 --- a/packages/orchestrator/pkg/sandbox/map_test.go +++ b/packages/orchestrator/pkg/sandbox/map_test.go @@ -4,6 +4,7 @@ package sandbox import ( "context" + "fmt" "sync" "testing" "testing/synctest" @@ -945,6 +946,144 @@ func TestSandboxCloseWithoutALiveEntryCountsNothing(t *testing.T) { require.NoError(t, live.Close(t.Context()), "Close must succeed with a live entry too") } +// Starts that have not reached MarkRunning hold a slot, so a node at its limit +// refuses the next create even while nothing is live yet. +func TestMapReserveWithinCountsStartsInFlight(t *testing.T) { + t.Parallel() + + sandboxes := NewSandboxesMap() + + a, err := sandboxes.ReserveWithin("sandbox-a", 2) + require.NoError(t, err) + b, err := sandboxes.ReserveWithin("sandbox-b", 2) + require.NoError(t, err) + require.Zero(t, sandboxes.Count(), "nothing is live yet") + + _, err = sandboxes.ReserveWithin("sandbox-c", 2) + require.ErrorIs(t, err, ErrNodeAtCapacity) + + a.Release() + c, err := sandboxes.ReserveWithin("sandbox-c", 2) + require.NoError(t, err, "a released start frees its slot") + b.Release() + c.Release() +} + +// Between MarkRunning and Release a start is both live and reserved. It must +// count once, or a node would refuse creates it has room for. +func TestMapReserveWithinCountsALiveReservationOnce(t *testing.T) { + t.Parallel() + + sandboxes := NewSandboxesMap() + + r, err := sandboxes.ReserveWithin("sandbox-1", 2) + require.NoError(t, err) + require.NoError(t, r.MarkRunning(t.Context(), testMapSandbox(t, "lifecycle-1"))) + + other, err := sandboxes.ReserveWithin("sandbox-2", 2) + require.NoError(t, err) + + _, err = sandboxes.ReserveWithin("sandbox-3", 2) + require.ErrorIs(t, err, ErrNodeAtCapacity) + + r.Release() + other.Release() + third, err := sandboxes.ReserveWithin("sandbox-3", 2) + require.NoError(t, err, "releasing the unstarted reservation frees its slot") + _, err = sandboxes.ReserveWithin("sandbox-4", 2) + require.ErrorIs(t, err, ErrNodeAtCapacity, "the live sandbox still holds its slot after the start finishes") + third.Release() +} + +// A checkpoint hold replaces the live entry while the old VM still runs, so it +// keeps holding the slot. +func TestMapReserveWithinCountsCheckpointHolds(t *testing.T) { + t.Parallel() + + sandboxes := NewSandboxesMap() + sbx := testMapSandbox(t, "lifecycle-1") + require.NoError(t, sandboxes.MarkRunning(t.Context(), sbx)) + + hold, err := sandboxes.MarkStoppingReserved(t.Context(), sbx.Runtime.SandboxID, sbx.LifecycleID) + require.NoError(t, err) + require.Zero(t, sandboxes.Count()) + + _, err = sandboxes.ReserveWithin("sandbox-2", 1) + require.ErrorIs(t, err, ErrNodeAtCapacity) + + hold.Release() + r, err := sandboxes.ReserveWithin("sandbox-2", 1) + require.NoError(t, err) + r.Release() +} + +// A duplicate ID is refused as a duplicate even on a full node, so the caller +// does not retry it elsewhere and create a second VM. +func TestMapReserveWithinRefusesDuplicatesBeforeCapacity(t *testing.T) { + t.Parallel() + + sandboxes := NewSandboxesMap() + r, err := sandboxes.ReserveWithin("sandbox-1", 1) + require.NoError(t, err) + t.Cleanup(r.Release) + + _, err = sandboxes.ReserveWithin("sandbox-1", 1) + require.ErrorIs(t, err, ErrSandboxAlreadyRunning) + require.NotErrorIs(t, err, ErrNodeAtCapacity) +} + +func TestMapReserveWithinRefusesNonPositiveLimits(t *testing.T) { + t.Parallel() + + sandboxes := NewSandboxesMap() + for _, limit := range []int64{0, -1} { + _, err := sandboxes.ReserveWithin("sandbox-1", limit) + require.ErrorIs(t, err, ErrNodeAtCapacity, "limit %d", limit) + } +} + +// Many concurrent creates for distinct IDs admit exactly limit of them. +func TestMapReserveWithinAdmitsExactlyTheLimitUnderConcurrency(t *testing.T) { + t.Parallel() + + const ( + limit = 5 + creators = 64 + ) + sandboxes := NewSandboxesMap() + + var ( + wg sync.WaitGroup + mu sync.Mutex + admitted []*Reservation + refused int + ) + start := make(chan struct{}) + for i := range creators { + wg.Go(func() { + <-start + r, err := sandboxes.ReserveWithin(fmt.Sprintf("sandbox-%d", i), limit) + mu.Lock() + defer mu.Unlock() + if err != nil { + assert.ErrorIs(t, err, ErrNodeAtCapacity) + refused++ + + return + } + admitted = append(admitted, r) + }) + } + close(start) + wg.Wait() + + require.Len(t, admitted, limit) + require.Equal(t, creators-limit, refused) + for _, r := range admitted { + r.Release() + } +} + func testMapSandbox(t *testing.T, lifecycleID string) *Sandbox { t.Helper() diff --git a/packages/orchestrator/pkg/server/create_duplicate_test.go b/packages/orchestrator/pkg/server/create_duplicate_test.go index 606f6f4504..45a17075ae 100644 --- a/packages/orchestrator/pkg/server/create_duplicate_test.go +++ b/packages/orchestrator/pkg/server/create_duplicate_test.go @@ -84,10 +84,12 @@ func TestCreate_ReleasesTheIDWhenItFails(t *testing.T) { t.Parallel() s := duplicateCreateTestServer() - s.info.MaxSandboxes.Store(0) // every create fails right after reserving + // Hold the only starting slot so the create fails right after reserving. + require.True(t, s.startingSandboxes.TryAcquire(1)) + t.Cleanup(func() { s.startingSandboxes.Release(1) }) _, err := s.Create(t.Context(), &orchestrator.SandboxCreateRequest{ - Sandbox: &orchestrator.SandboxConfig{SandboxId: "sandbox-1", Snapshot: true}, + Sandbox: &orchestrator.SandboxConfig{SandboxId: "sandbox-1"}, }) st, ok := status.FromError(err) require.True(t, ok) @@ -98,6 +100,35 @@ func TestCreate_ReleasesTheIDWhenItFails(t *testing.T) { r.Release() } +// A start still in flight holds a slot against MaxSandboxes. With one slot +// taken by an in-flight create, the next create is refused before any +// resource is touched: the template cache is nil here, so reaching it would +// panic. +func TestCreate_CountsStartsInFlightAgainstTheNodeLimit(t *testing.T) { + t.Parallel() + + s := duplicateCreateTestServer() + s.info.MaxSandboxes.Store(1) + + inFlight, err := s.sandboxFactory.Sandboxes.Reserve("sandbox-in-flight") + require.NoError(t, err) + + for _, snapshot := range []bool{false, true} { + _, err = s.Create(t.Context(), &orchestrator.SandboxCreateRequest{ + Sandbox: &orchestrator.SandboxConfig{SandboxId: "sandbox-1", Snapshot: snapshot}, + }) + require.Equal(t, codes.ResourceExhausted, status.Code(err), "snapshot=%t", snapshot) + } + assert.Zero(t, s.sandboxFactory.Sandboxes.Count(), "nothing was registered") + + // The refusal took no hold of its own; the in-flight create keeps its slot + // until it ends, and then the ID is free. + inFlight.Release() + r, err := s.sandboxFactory.Sandboxes.ReserveWithin("sandbox-1", 1) + require.NoError(t, err) + r.Release() +} + func TestFailedStartReturnsBeforeCleanupAndRetainsReservation(t *testing.T) { t.Parallel() diff --git a/packages/orchestrator/pkg/server/sandboxes.go b/packages/orchestrator/pkg/server/sandboxes.go index a667e037d2..df76db0750 100644 --- a/packages/orchestrator/pkg/server/sandboxes.go +++ b/packages/orchestrator/pkg/server/sandboxes.go @@ -210,7 +210,16 @@ func (s *Server) Create(ctx context.Context, req *orchestrator.SandboxCreateRequ } } - reservation, err := s.sandboxFactory.Sandboxes.Reserve(req.GetSandbox().GetSandboxId()) + // The limit counts starts still in flight, not only live sandboxes, and is + // checked atomically with taking the ID. Otherwise concurrent creates could + // all see the same free slot and push the node past its limit. + maxRunningSandboxesPerNode := s.info.MaxSandboxes.Load() + reservation, err := s.sandboxFactory.Sandboxes.ReserveWithin(req.GetSandbox().GetSandboxId(), maxRunningSandboxesPerNode) + if errors.Is(err, sandbox.ErrNodeAtCapacity) { + telemetry.ReportEvent(ctx, "max number of running sandboxes reached") + + return nil, status.Errorf(codes.ResourceExhausted, "max number of running sandboxes on node reached (%d), please retry", maxRunningSandboxesPerNode) + } if err != nil { return nil, s.sandboxAlreadyRunning(ctx, req.GetSandbox().GetSandboxId(), req.GetSandbox().GetExecutionId(), err) } @@ -219,15 +228,6 @@ func (s *Server) Create(ctx context.Context, req *orchestrator.SandboxCreateRequ s.finishSandboxStart(ctx, reservation, rollback, createErr) }() - maxRunningSandboxesPerNode := s.info.MaxSandboxes.Load() - - runningSandboxes := int64(s.sandboxFactory.Sandboxes.Count()) - if runningSandboxes >= maxRunningSandboxesPerNode { - telemetry.ReportEvent(ctx, "max number of running sandboxes reached") - - return nil, status.Errorf(codes.ResourceExhausted, "max number of running sandboxes on node reached (%d), please retry", maxRunningSandboxesPerNode) - } - // Check if we've reached the max number of starting instances on this node if req.GetSandbox().GetSnapshot() { err := s.waitForAcquire(ctx)