Skip to content
Open
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
40 changes: 40 additions & 0 deletions packages/orchestrator/pkg/sandbox/map.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()

Expand All @@ -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()
Expand Down
139 changes: 139 additions & 0 deletions packages/orchestrator/pkg/sandbox/map_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ package sandbox

import (
"context"
"fmt"
"sync"
"testing"
"testing/synctest"
Expand Down Expand Up @@ -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()

Expand Down
35 changes: 33 additions & 2 deletions packages/orchestrator/pkg/server/create_duplicate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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()

Expand Down
20 changes: 10 additions & 10 deletions packages/orchestrator/pkg/server/sandboxes.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand All @@ -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)
Expand Down