From 567ed7ac27ea3f8ab1d96d0a08e0d358b96cc0e4 Mon Sep 17 00:00:00 2001 From: Luke Lombardi <33990301+luke-lombardi@users.noreply.github.com> Date: Tue, 28 Jul 2026 16:58:52 -0400 Subject: [PATCH 1/5] feat: route single-GPU workers to dedicated nodes --- pkg/scheduler/pool.go | 22 ++++ pkg/scheduler/pool_job_spec_test.go | 185 ++++++++++++++++++++++++++++ pkg/scheduler/pool_local.go | 2 +- pkg/scheduler/pool_provider.go | 37 +++--- pkg/types/config.go | 5 +- 5 files changed, 234 insertions(+), 17 deletions(-) create mode 100644 pkg/scheduler/pool_job_spec_test.go diff --git a/pkg/scheduler/pool.go b/pkg/scheduler/pool.go index 4d8fb31c3..f80a075e4 100644 --- a/pkg/scheduler/pool.go +++ b/pkg/scheduler/pool.go @@ -76,6 +76,28 @@ type WorkerPoolCapacity struct { PendingGpu uint } +func workerNodeSelector(jobSpec types.WorkerPoolJobSpecConfig, gpuCount uint32) map[string]string { + var selector map[string]string + if len(jobSpec.NodeSelector) > 0 { + selector = make(map[string]string, len(jobSpec.NodeSelector)) + for key, value := range jobSpec.NodeSelector { + selector[key] = value + } + } + + if gpuCount != 1 || len(jobSpec.SingleGPUNodeSelector) == 0 { + return selector + } + + if selector == nil { + selector = make(map[string]string, len(jobSpec.SingleGPUNodeSelector)) + } + for key, value := range jobSpec.SingleGPUNodeSelector { + selector[key] = value + } + return selector +} + type WorkerPoolControllerOptions struct { Name string Context context.Context diff --git a/pkg/scheduler/pool_job_spec_test.go b/pkg/scheduler/pool_job_spec_test.go new file mode 100644 index 000000000..83dd8c537 --- /dev/null +++ b/pkg/scheduler/pool_job_spec_test.go @@ -0,0 +1,185 @@ +package scheduler + +import ( + "testing" + + "github.com/beam-cloud/beta9/pkg/types" + "github.com/tj/assert" +) + +func TestWorkerNodeSelector(t *testing.T) { + jobSpec := types.WorkerPoolJobSpecConfig{ + NodeSelector: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + SingleGPUNodeSelector: map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + "karpenter.k8s.aws/instance-gpu-count": "1", + }, + } + + tests := []struct { + name string + gpuCount uint32 + want map[string]string + }{ + { + name: "CPU uses base selector", + want: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + }, + { + name: "single GPU overlays selector", + gpuCount: 1, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + "kubernetes.io/arch": "amd64", + "karpenter.k8s.aws/instance-gpu-count": "1", + }, + }, + { + name: "multiple GPUs use base selector", + gpuCount: 2, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + }, + { + name: "four GPUs use base selector", + gpuCount: 4, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got := workerNodeSelector(jobSpec, test.gpuCount) + assert.Equal(t, test.want, got) + + got["mutated"] = "true" + assert.Equal(t, "gpu", jobSpec.NodeSelector["karpenter.sh/nodepool"]) + assert.Equal(t, "gpu-single", jobSpec.SingleGPUNodeSelector["karpenter.sh/nodepool"]) + _, baseMutated := jobSpec.NodeSelector["mutated"] + _, singleMutated := jobSpec.SingleGPUNodeSelector["mutated"] + assert.False(t, baseMutated) + assert.False(t, singleMutated) + }) + } +} + +func TestWorkerNodeSelectorSingleGPUWithoutBase(t *testing.T) { + jobSpec := types.WorkerPoolJobSpecConfig{ + SingleGPUNodeSelector: map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + }, + } + + assert.Equal(t, map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + }, workerNodeSelector(jobSpec, 1)) + assert.Nil(t, workerNodeSelector(jobSpec, 0)) +} + +func TestWorkerJobBuildersUseGPUCountSelectorsWithoutMutatingConfig(t *testing.T) { + jobSpec := types.WorkerPoolJobSpecConfig{ + NodeSelector: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + SingleGPUNodeSelector: map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + "karpenter.k8s.aws/instance-gpu-count": "1", + }, + } + workerConfig := types.WorkerConfig{ + DefaultWorkerCPURequest: 1000, + DefaultWorkerMemoryRequest: 1024, + } + localController := &LocalKubernetesWorkerPoolController{ + name: "local-gpu", + config: types.AppConfig{Worker: workerConfig}, + workerPoolConfig: types.WorkerPoolConfig{JobSpec: jobSpec}, + } + providerController := &ProviderWorkerPoolController{ + name: "provider-gpu", + config: types.AppConfig{Worker: workerConfig}, + workerPoolConfig: types.WorkerPoolConfig{JobSpec: jobSpec}, + } + + builders := []struct { + name string + build func(uint32) map[string]string + }{ + { + name: "local Job", + build: func(gpuCount uint32) map[string]string { + job, _ := localController.createWorkerJob("worker-local", 1000, 1024, "A10G", gpuCount, "token") + return job.Spec.Template.Spec.NodeSelector + }, + }, + { + name: "provider Job", + build: func(gpuCount uint32) map[string]string { + job, _ := providerController.buildWorkerJob("worker-provider", "machine-provider", 1000, 1024, "A10G", gpuCount, nil) + return job.Spec.Template.Spec.NodeSelector + }, + }, + } + tests := []struct { + name string + gpuCount uint32 + want map[string]string + }{ + { + name: "one GPU overlays selector", + gpuCount: 1, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu-single", + "kubernetes.io/arch": "amd64", + "karpenter.k8s.aws/instance-gpu-count": "1", + }, + }, + { + name: "two GPUs retain base selector", + gpuCount: 2, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + }, + { + name: "four GPUs retain base selector", + gpuCount: 4, + want: map[string]string{ + "karpenter.sh/nodepool": "gpu", + "kubernetes.io/arch": "amd64", + }, + }, + } + + for _, builder := range builders { + t.Run(builder.name, func(t *testing.T) { + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + selector := builder.build(test.gpuCount) + assert.Equal(t, test.want, selector) + + selector["mutated"] = "true" + assert.Equal(t, "gpu", jobSpec.NodeSelector["karpenter.sh/nodepool"]) + assert.Equal(t, "gpu-single", jobSpec.SingleGPUNodeSelector["karpenter.sh/nodepool"]) + _, baseMutated := jobSpec.NodeSelector["mutated"] + _, singleMutated := jobSpec.SingleGPUNodeSelector["mutated"] + assert.False(t, baseMutated) + assert.False(t, singleMutated) + }) + } + }) + } +} diff --git a/pkg/scheduler/pool_local.go b/pkg/scheduler/pool_local.go index 925229f6f..d96cd3763 100644 --- a/pkg/scheduler/pool_local.go +++ b/pkg/scheduler/pool_local.go @@ -253,7 +253,7 @@ func (wpc *LocalKubernetesWorkerPoolController) createWorkerJob(workerId string, HostNetwork: wpc.config.Worker.HostNetwork, ImagePullSecrets: imagePullSecrets, RestartPolicy: corev1.RestartPolicyOnFailure, - NodeSelector: wpc.workerPoolConfig.JobSpec.NodeSelector, + NodeSelector: workerNodeSelector(wpc.workerPoolConfig.JobSpec, workerGpuCount), Containers: containers, Volumes: wpc.getWorkerVolumes(workerMemory), EnableServiceLinks: ptr.To(false), diff --git a/pkg/scheduler/pool_provider.go b/pkg/scheduler/pool_provider.go index 493fa75d6..a2e3d0e78 100644 --- a/pkg/scheduler/pool_provider.go +++ b/pkg/scheduler/pool_provider.go @@ -329,6 +329,25 @@ func (wpc *ProviderWorkerPoolController) createWorkerOnMachine(workerId, machine } func (wpc *ProviderWorkerPoolController) createWorkerJob(workerId, machineId string, cpu int64, memory int64, gpuType string, gpuCount uint32, token string) (*batchv1.Job, *types.Worker, error) { + workerGpuType := wpc.workerPoolConfig.GPUType + env, err := wpc.getWorkerEnvironment(workerId, machineId, cpu, memory, workerGpuType, gpuCount, token) + if err != nil { + return nil, nil, err + } + + job, worker := wpc.buildWorkerJob(workerId, machineId, cpu, memory, workerGpuType, gpuCount, env) + return job, worker, nil +} + +func (wpc *ProviderWorkerPoolController) buildWorkerJob( + workerId string, + machineId string, + workerCpu int64, + workerMemory int64, + workerGpuType string, + workerGpuCount uint32, + env []corev1.EnvVar, +) (*batchv1.Job, *types.Worker) { jobName := fmt.Sprintf("%s-%s-%s", Beta9WorkerJobPrefix, wpc.name, workerId) labels := map[string]string{ "app": Beta9WorkerLabelValue, @@ -339,11 +358,6 @@ func (wpc *ProviderWorkerPoolController) createWorkerJob(workerId, machineId str PrometheusScrapeKey: strconv.FormatBool(wpc.config.Monitoring.Prometheus.ScrapeWorkers), } - workerCpu := cpu - workerMemory := memory - workerGpuType := wpc.workerPoolConfig.GPUType - workerGpuCount := gpuCount - workerImage := fmt.Sprintf("%s/%s:%s", wpc.config.Worker.ImageRegistry, wpc.config.Worker.ImageName, @@ -353,18 +367,13 @@ func (wpc *ProviderWorkerPoolController) createWorkerJob(workerId, machineId str resources := corev1.ResourceRequirements{} if workerGpuType != "" { resources.Requests = corev1.ResourceList{ - "nvidia.com/gpu": *resource.NewQuantity(int64(gpuCount), resource.DecimalSI), + "nvidia.com/gpu": *resource.NewQuantity(int64(workerGpuCount), resource.DecimalSI), } resources.Limits = corev1.ResourceList{ - "nvidia.com/gpu": *resource.NewQuantity(int64(gpuCount), resource.DecimalSI), + "nvidia.com/gpu": *resource.NewQuantity(int64(workerGpuCount), resource.DecimalSI), } } - env, err := wpc.getWorkerEnvironment(workerId, machineId, workerCpu, workerMemory, workerGpuType, workerGpuCount, token) - if err != nil { - return nil, nil, err - } - containers := []corev1.Container{ { Name: defaultContainerName, @@ -393,7 +402,7 @@ func (wpc *ProviderWorkerPoolController) createWorkerJob(workerId, machineId str HostNetwork: true, ImagePullSecrets: imagePullSecrets, RestartPolicy: corev1.RestartPolicyOnFailure, - NodeSelector: wpc.workerPoolConfig.JobSpec.NodeSelector, + NodeSelector: workerNodeSelector(wpc.workerPoolConfig.JobSpec, workerGpuCount), Containers: containers, Volumes: wpc.getWorkerVolumes(workerMemory), EnableServiceLinks: ptr.To(false), @@ -433,7 +442,7 @@ func (wpc *ProviderWorkerPoolController) createWorkerJob(workerId, machineId str Runtime: wpc.workerPoolConfig.ContainerRuntime, BuildVersion: wpc.config.Worker.ImageTag, Preemptable: wpc.workerPoolConfig.Preemptable, - }, nil + } } func (wpc *ProviderWorkerPoolController) getWorkerEnvironment(workerId, machineId string, cpu int64, memory int64, gpuType string, gpuCount uint32, token string) ([]corev1.EnvVar, error) { diff --git a/pkg/types/config.go b/pkg/types/config.go index 53a4adcc2..da2959202 100644 --- a/pkg/types/config.go +++ b/pkg/types/config.go @@ -534,8 +534,9 @@ func (c RuntimeConfig) WithDefaults(runtime string) RuntimeConfig { } type WorkerPoolJobSpecConfig struct { - NodeSelector map[string]string `key:"nodeSelector" json:"node_selector"` - Env []corev1.EnvVar `key:"env" json:"env"` + NodeSelector map[string]string `key:"nodeSelector" json:"node_selector"` + SingleGPUNodeSelector map[string]string `key:"singleGpuNodeSelector" json:"single_gpu_node_selector"` + Env []corev1.EnvVar `key:"env" json:"env"` // Mimics corev1.Volume since that type doesn't currently serialize correctly Volumes []struct { From 14522a6ad886c91a22fb8a9745cbe03b43692223 Mon Sep 17 00:00:00 2001 From: Luke Lombardi <33990301+luke-lombardi@users.noreply.github.com> Date: Tue, 28 Jul 2026 16:59:07 -0400 Subject: [PATCH 2/5] feat: pack work onto best-fit workers --- pkg/scheduler/scheduler.go | 7 +- pkg/scheduler/scheduler_test.go | 143 ++++++++++++++++++++++++++++---- 2 files changed, 134 insertions(+), 16 deletions(-) diff --git a/pkg/scheduler/scheduler.go b/pkg/scheduler/scheduler.go index dc0ead629..2f111b33b 100644 --- a/pkg/scheduler/scheduler.go +++ b/pkg/scheduler/scheduler.go @@ -1164,7 +1164,12 @@ func (s *Scheduler) selectWorkerFromWorkersByStatus(workers []*types.Worker, req if scoredWorkers[i].score != scoredWorkers[j].score { return scoredWorkers[i].score > scoredWorkers[j].score } - return workerFreeCapacityScore(scoredWorkers[i].worker, request) > workerFreeCapacityScore(scoredWorkers[j].worker, request) + iCapacity := workerFreeCapacityScore(scoredWorkers[i].worker, request) + jCapacity := workerFreeCapacityScore(scoredWorkers[j].worker, request) + if iCapacity != jCapacity { + return iCapacity < jCapacity + } + return scoredWorkers[i].worker.Id < scoredWorkers[j].worker.Id }) return scoredWorkers[0].worker, nil diff --git a/pkg/scheduler/scheduler_test.go b/pkg/scheduler/scheduler_test.go index ba7464372..c336f1960 100644 --- a/pkg/scheduler/scheduler_test.go +++ b/pkg/scheduler/scheduler_test.go @@ -1495,12 +1495,12 @@ func TestProcessRequestBatchDoesNotOverScheduleWorkerSnapshot(t *testing.T) { assert.Equal(t, int64(2), updatedWorker.ResourceVersion) } -func TestProcessRequestBatchSpreadsAcrossEqualWorkers(t *testing.T) { +func TestProcessRequestBatchPacksOntoEqualWorkers(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) firstWorker := &types.Worker{ - Id: uuid.New().String(), + Id: "worker-a", Status: types.WorkerStatusAvailable, TotalCpu: 200, TotalMemory: 250, @@ -1512,7 +1512,7 @@ func TestProcessRequestBatchSpreadsAcrossEqualWorkers(t *testing.T) { assert.Nil(t, err) secondWorker := &types.Worker{ - Id: uuid.New().String(), + Id: "worker-b", Status: types.WorkerStatusAvailable, TotalCpu: 200, TotalMemory: 250, @@ -1544,17 +1544,21 @@ func TestProcessRequestBatchSpreadsAcrossEqualWorkers(t *testing.T) { wb.processRequestBatch(requests, workers) - firstQueued, err := wb.workerRepo.GetNextContainerRequest(firstWorker.Id) + for i := 0; i < len(requests); i++ { + queued, err := wb.workerRepo.GetNextContainerRequest(firstWorker.Id) + assert.Nil(t, err) + assert.NotNil(t, queued) + } + firstExtra, err := wb.workerRepo.GetNextContainerRequest(firstWorker.Id) assert.Nil(t, err) + assert.Nil(t, firstExtra) + secondQueued, err := wb.workerRepo.GetNextContainerRequest(secondWorker.Id) assert.Nil(t, err) - - assert.NotNil(t, firstQueued) - assert.NotNil(t, secondQueued) - assert.NotEqual(t, firstQueued.ContainerId, secondQueued.ContainerId) + assert.Nil(t, secondQueued) } -func TestProcessRequestBatchSchedulesTinySandboxBurstByRequestedCapacity(t *testing.T) { +func TestProcessRequestBatchPacksTinySandboxBurstByRequestedCapacity(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) @@ -1565,7 +1569,7 @@ func TestProcessRequestBatchSchedulesTinySandboxBurstByRequestedCapacity(t *test workers := []*types.Worker{ { - Id: uuid.New().String(), + Id: "worker-a", Status: types.WorkerStatusAvailable, TotalCpu: 16000, TotalMemory: 32000, @@ -1575,7 +1579,7 @@ func TestProcessRequestBatchSchedulesTinySandboxBurstByRequestedCapacity(t *test Runtime: types.ContainerRuntimeRunc.String(), }, { - Id: uuid.New().String(), + Id: "worker-b", Status: types.WorkerStatusAvailable, TotalCpu: 16000, TotalMemory: 32000, @@ -1608,7 +1612,7 @@ func TestProcessRequestBatchSchedulesTinySandboxBurstByRequestedCapacity(t *test assert.Nil(t, err) wb.processRequestBatch(requests, workerSnapshot) - for _, worker := range workers { + for index, worker := range workers { queued := 0 for { request, err := wb.workerRepo.GetNextContainerRequest(worker.Id) @@ -1618,7 +1622,11 @@ func TestProcessRequestBatchSchedulesTinySandboxBurstByRequestedCapacity(t *test } queued++ } - assert.Equal(t, 35, queued) + if index == 0 { + assert.Equal(t, len(requests), queued) + } else { + assert.Equal(t, 0, queued) + } } } @@ -2824,8 +2832,8 @@ func TestSelectGPUWorkerDoesNotMutatePriority(t *testing.T) { t4Worker := &types.Worker{ Id: uuid.New().String(), Status: types.WorkerStatusAvailable, - FreeCpu: 1000, - FreeMemory: 1250, + FreeCpu: 2000, + FreeMemory: 2500, FreeGpuCount: 1, Gpu: "T4", Priority: 7, @@ -2844,6 +2852,111 @@ func TestSelectGPUWorkerDoesNotMutatePriority(t *testing.T) { assert.Equal(t, int32(7), t4Worker.Priority) } +func TestSelectWorkerBestFitTieBreak(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + + request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} + roomyWorker := &types.Worker{ + Id: "roomy", + Status: types.WorkerStatusAvailable, + FreeCpu: 4000, + FreeMemory: 8000, + Priority: 7, + } + bestFitWorker := &types.Worker{ + Id: "best-fit", + Status: types.WorkerStatusAvailable, + FreeCpu: 1500, + FreeMemory: 2000, + Priority: 7, + } + + worker, err := wb.selectWorkerFromWorkers([]*types.Worker{roomyWorker, bestFitWorker}, request) + assert.Nil(t, err) + assert.Equal(t, bestFitWorker.Id, worker.Id) +} + +func TestSelectWorkerBestFitPreservesPriority(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + + request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} + highPriorityWorker := &types.Worker{ + Id: "high-priority", + Status: types.WorkerStatusAvailable, + FreeCpu: 4000, + FreeMemory: 8000, + Priority: 8, + } + bestFitWorker := &types.Worker{ + Id: "best-fit", + Status: types.WorkerStatusAvailable, + FreeCpu: 1000, + FreeMemory: 1000, + Priority: 7, + } + + worker, err := wb.selectWorkerFromWorkers([]*types.Worker{bestFitWorker, highPriorityWorker}, request) + assert.Nil(t, err) + assert.Equal(t, highPriorityWorker.Id, worker.Id) +} + +func TestSelectWorkerBestFitPreservesAvailability(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + + request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} + availableWorker := &types.Worker{ + Id: "available", + Status: types.WorkerStatusAvailable, + FreeCpu: 4000, + FreeMemory: 8000, + Priority: 7, + } + pendingBestFitWorker := &types.Worker{ + Id: "pending-best-fit", + Status: types.WorkerStatusPending, + FreeCpu: 1000, + FreeMemory: 1000, + Priority: 7, + } + + worker, err := wb.selectWorkerFromWorkersByStatus( + []*types.Worker{pendingBestFitWorker, availableWorker}, + request, + types.WorkerStatusAvailable, + types.WorkerStatusPending, + ) + assert.Nil(t, err) + assert.Equal(t, availableWorker.Id, worker.Id) +} + +func TestSelectWorkerBestFitUsesWorkerIDForDeterministicTies(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + + request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} + workerB := &types.Worker{ + Id: "worker-b", + Status: types.WorkerStatusAvailable, + FreeCpu: 2000, + FreeMemory: 2000, + Priority: 7, + } + workerA := &types.Worker{ + Id: "worker-a", + Status: types.WorkerStatusAvailable, + FreeCpu: 2000, + FreeMemory: 2000, + Priority: 7, + } + + worker, err := wb.selectWorkerFromWorkers([]*types.Worker{workerB, workerA}, request) + assert.Nil(t, err) + assert.Equal(t, workerA.Id, worker.Id) +} + func TestRequiresPoolSelectorWorker(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) From 2b9886404ec89758f6b23a7a76c377ee7c4238cf Mon Sep 17 00:00:00 2001 From: Luke Lombardi <33990301+luke-lombardi@users.noreply.github.com> Date: Wed, 29 Jul 2026 09:09:58 -0400 Subject: [PATCH 3/5] style: gofmt worker pool config --- pkg/types/config.go | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pkg/types/config.go b/pkg/types/config.go index da2959202..629cbe26d 100644 --- a/pkg/types/config.go +++ b/pkg/types/config.go @@ -176,12 +176,12 @@ type ContainerCostHookConfig struct { } type GatewayServiceConfig struct { - Host string `key:"host" json:"host"` - InvokeURLType string `key:"invokeURLType" json:"invoke_url_type"` - GRPC GRPCConfig `key:"grpc" json:"grpc"` - HTTP HTTPConfig `key:"http" json:"http"` - ShutdownTimeout time.Duration `key:"shutdownTimeout" json:"shutdown_timeout"` - StubLimits StubLimits `key:"stubLimits" json:"stub_limits"` + Host string `key:"host" json:"host"` + InvokeURLType string `key:"invokeURLType" json:"invoke_url_type"` + GRPC GRPCConfig `key:"grpc" json:"grpc"` + HTTP HTTPConfig `key:"http" json:"http"` + ShutdownTimeout time.Duration `key:"shutdownTimeout" json:"shutdown_timeout"` + StubLimits StubLimits `key:"stubLimits" json:"stub_limits"` } type FileServiceConfig struct { From 11391eba5e17a579032c4b446023e315dd2b60a2 Mon Sep 17 00:00:00 2001 From: Luke Lombardi <33990301+luke-lombardi@users.noreply.github.com> Date: Wed, 29 Jul 2026 13:04:28 -0400 Subject: [PATCH 4/5] feat: gate best-fit worker packing --- pkg/scheduler/scheduler.go | 3 + pkg/scheduler/scheduler_test.go | 345 +++++++++++++++++++------------- pkg/types/config.go | 1 + 3 files changed, 213 insertions(+), 136 deletions(-) diff --git a/pkg/scheduler/scheduler.go b/pkg/scheduler/scheduler.go index 2f111b33b..7163f38d1 100644 --- a/pkg/scheduler/scheduler.go +++ b/pkg/scheduler/scheduler.go @@ -1166,6 +1166,9 @@ func (s *Scheduler) selectWorkerFromWorkersByStatus(workers []*types.Worker, req } iCapacity := workerFreeCapacityScore(scoredWorkers[i].worker, request) jCapacity := workerFreeCapacityScore(scoredWorkers[j].worker, request) + if !s.config.Worker.PreferBestFit { + return iCapacity > jCapacity + } if iCapacity != jCapacity { return iCapacity < jCapacity } diff --git a/pkg/scheduler/scheduler_test.go b/pkg/scheduler/scheduler_test.go index c336f1960..ad2534787 100644 --- a/pkg/scheduler/scheduler_test.go +++ b/pkg/scheduler/scheduler_test.go @@ -1495,138 +1495,184 @@ func TestProcessRequestBatchDoesNotOverScheduleWorkerSnapshot(t *testing.T) { assert.Equal(t, int64(2), updatedWorker.ResourceVersion) } -func TestProcessRequestBatchPacksOntoEqualWorkers(t *testing.T) { - wb, err := NewSchedulerForTest() - assert.Nil(t, err) - - firstWorker := &types.Worker{ - Id: "worker-a", - Status: types.WorkerStatusAvailable, - TotalCpu: 200, - TotalMemory: 250, - FreeCpu: 200, - FreeMemory: 250, - PoolName: "beta9-cpu", - } - err = wb.workerRepo.AddWorker(firstWorker) - assert.Nil(t, err) - - secondWorker := &types.Worker{ - Id: "worker-b", - Status: types.WorkerStatusAvailable, - TotalCpu: 200, - TotalMemory: 250, - FreeCpu: 200, - FreeMemory: 250, - PoolName: "beta9-cpu", - } - err = wb.workerRepo.AddWorker(secondWorker) - assert.Nil(t, err) - - workers, err := wb.workerRepo.GetAllWorkers() - assert.Nil(t, err) - - requests := []*types.ContainerRequest{ +func TestProcessRequestBatchHonorsPackingPreference(t *testing.T) { + tests := []struct { + name string + enableBestFit bool + expectedFirst int + expectedSecond int + }{ { - ContainerId: uuid.New().String(), - Cpu: 100, - Memory: 100, - Timestamp: time.Now(), + name: "spread by default", + expectedFirst: 1, + expectedSecond: 1, }, { - ContainerId: uuid.New().String(), - Cpu: 100, - Memory: 100, - Timestamp: time.Now(), + name: "best-fit when enabled", + enableBestFit: true, + expectedFirst: 2, + expectedSecond: 0, }, } - setPendingSchedulerRequests(t, wb, requests...) - - wb.processRequestBatch(requests, workers) - for i := 0; i < len(requests); i++ { - queued, err := wb.workerRepo.GetNextContainerRequest(firstWorker.Id) - assert.Nil(t, err) - assert.NotNil(t, queued) - } - firstExtra, err := wb.workerRepo.GetNextContainerRequest(firstWorker.Id) - assert.Nil(t, err) - assert.Nil(t, firstExtra) - - secondQueued, err := wb.workerRepo.GetNextContainerRequest(secondWorker.Id) - assert.Nil(t, err) - assert.Nil(t, secondQueued) -} + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + assert.False(t, wb.config.Worker.PreferBestFit) + if test.enableBestFit { + wb.config.Worker.PreferBestFit = true + } -func TestProcessRequestBatchPacksTinySandboxBurstByRequestedCapacity(t *testing.T) { - wb, err := NewSchedulerForTest() - assert.Nil(t, err) + workers := []*types.Worker{ + { + Id: "worker-a", + Status: types.WorkerStatusAvailable, + TotalCpu: 200, + TotalMemory: 250, + FreeCpu: 200, + FreeMemory: 250, + PoolName: "beta9-cpu", + }, + { + Id: "worker-b", + Status: types.WorkerStatusAvailable, + TotalCpu: 200, + TotalMemory: 250, + FreeCpu: 200, + FreeMemory: 250, + PoolName: "beta9-cpu", + }, + } + for _, worker := range workers { + err = wb.workerRepo.AddWorker(worker) + assert.Nil(t, err) + } - wb.config.Worker.Pools["beta9-cpu"] = types.WorkerPoolConfig{ - ContainerRuntime: types.ContainerRuntimeRunc.String(), - ContainerStartConcurrency: 32, + workerSnapshot, err := wb.workerRepo.GetAllWorkers() + assert.Nil(t, err) + requests := []*types.ContainerRequest{ + { + ContainerId: uuid.New().String(), + Cpu: 100, + Memory: 100, + Timestamp: time.Now(), + }, + { + ContainerId: uuid.New().String(), + Cpu: 100, + Memory: 100, + Timestamp: time.Now(), + }, + } + setPendingSchedulerRequests(t, wb, requests...) + wb.processRequestBatch(requests, workerSnapshot) + + queuedCounts := make([]int, len(workers)) + for index, worker := range workers { + for { + request, err := wb.workerRepo.GetNextContainerRequest(worker.Id) + assert.Nil(t, err) + if request == nil { + break + } + queuedCounts[index]++ + } + } + assert.Equal(t, test.expectedFirst, queuedCounts[0]) + assert.Equal(t, test.expectedSecond, queuedCounts[1]) + }) } +} - workers := []*types.Worker{ +func TestProcessRequestBatchHonorsPackingPreferenceForTinySandboxBurst(t *testing.T) { + tests := []struct { + name string + enableBestFit bool + expectedPerWorker []int + }{ { - Id: "worker-a", - Status: types.WorkerStatusAvailable, - TotalCpu: 16000, - TotalMemory: 32000, - FreeCpu: 16000, - FreeMemory: 32000, - PoolName: "beta9-cpu", - Runtime: types.ContainerRuntimeRunc.String(), + name: "spread by default", + expectedPerWorker: []int{35, 35}, }, { - Id: "worker-b", - Status: types.WorkerStatusAvailable, - TotalCpu: 16000, - TotalMemory: 32000, - FreeCpu: 16000, - FreeMemory: 32000, - PoolName: "beta9-cpu", - Runtime: types.ContainerRuntimeRunc.String(), + name: "best-fit when enabled", + enableBestFit: true, + expectedPerWorker: []int{70, 0}, }, } - for _, worker := range workers { - err = wb.workerRepo.AddWorker(worker) - assert.Nil(t, err) - } - requests := make([]*types.ContainerRequest, 70) - for i := range requests { - requests[i] = &types.ContainerRequest{ - ContainerId: uuid.New().String(), - Cpu: 100, - Memory: 100, - Timestamp: time.Now(), - Stub: types.StubWithRelated{ - Stub: types.Stub{Type: types.StubType(types.StubTypeSandbox)}, - }, - } - } - setPendingSchedulerRequests(t, wb, requests...) + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + assert.False(t, wb.config.Worker.PreferBestFit) + if test.enableBestFit { + wb.config.Worker.PreferBestFit = true + } + wb.config.Worker.Pools["beta9-cpu"] = types.WorkerPoolConfig{ + ContainerRuntime: types.ContainerRuntimeRunc.String(), + ContainerStartConcurrency: 32, + } - workerSnapshot, err := wb.workerRepo.GetAllWorkers() - assert.Nil(t, err) - wb.processRequestBatch(requests, workerSnapshot) + workers := []*types.Worker{ + { + Id: "worker-a", + Status: types.WorkerStatusAvailable, + TotalCpu: 16000, + TotalMemory: 32000, + FreeCpu: 16000, + FreeMemory: 32000, + PoolName: "beta9-cpu", + Runtime: types.ContainerRuntimeRunc.String(), + }, + { + Id: "worker-b", + Status: types.WorkerStatusAvailable, + TotalCpu: 16000, + TotalMemory: 32000, + FreeCpu: 16000, + FreeMemory: 32000, + PoolName: "beta9-cpu", + Runtime: types.ContainerRuntimeRunc.String(), + }, + } + for _, worker := range workers { + err = wb.workerRepo.AddWorker(worker) + assert.Nil(t, err) + } - for index, worker := range workers { - queued := 0 - for { - request, err := wb.workerRepo.GetNextContainerRequest(worker.Id) + requests := make([]*types.ContainerRequest, 70) + for i := range requests { + requests[i] = &types.ContainerRequest{ + ContainerId: uuid.New().String(), + Cpu: 100, + Memory: 100, + Timestamp: time.Now(), + Stub: types.StubWithRelated{ + Stub: types.Stub{Type: types.StubType(types.StubTypeSandbox)}, + }, + } + } + setPendingSchedulerRequests(t, wb, requests...) + + workerSnapshot, err := wb.workerRepo.GetAllWorkers() assert.Nil(t, err) - if request == nil { - break + wb.processRequestBatch(requests, workerSnapshot) + + for index, worker := range workers { + queued := 0 + for { + request, err := wb.workerRepo.GetNextContainerRequest(worker.Id) + assert.Nil(t, err) + if request == nil { + break + } + queued++ + } + assert.Equal(t, test.expectedPerWorker[index], queued) } - queued++ - } - if index == 0 { - assert.Equal(t, len(requests), queued) - } else { - assert.Equal(t, 0, queued) - } + }) } } @@ -2852,34 +2898,59 @@ func TestSelectGPUWorkerDoesNotMutatePriority(t *testing.T) { assert.Equal(t, int32(7), t4Worker.Priority) } -func TestSelectWorkerBestFitTieBreak(t *testing.T) { - wb, err := NewSchedulerForTest() - assert.Nil(t, err) - - request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} - roomyWorker := &types.Worker{ - Id: "roomy", - Status: types.WorkerStatusAvailable, - FreeCpu: 4000, - FreeMemory: 8000, - Priority: 7, - } - bestFitWorker := &types.Worker{ - Id: "best-fit", - Status: types.WorkerStatusAvailable, - FreeCpu: 1500, - FreeMemory: 2000, - Priority: 7, +func TestSelectWorkerCapacityTieBreakHonorsPackingPreference(t *testing.T) { + tests := []struct { + name string + enableBestFit bool + expectedID string + }{ + { + name: "spread by default", + expectedID: "roomy", + }, + { + name: "best-fit when enabled", + enableBestFit: true, + expectedID: "best-fit", + }, } - worker, err := wb.selectWorkerFromWorkers([]*types.Worker{roomyWorker, bestFitWorker}, request) - assert.Nil(t, err) - assert.Equal(t, bestFitWorker.Id, worker.Id) + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + wb, err := NewSchedulerForTest() + assert.Nil(t, err) + assert.False(t, wb.config.Worker.PreferBestFit) + if test.enableBestFit { + wb.config.Worker.PreferBestFit = true + } + + request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} + roomyWorker := &types.Worker{ + Id: "roomy", + Status: types.WorkerStatusAvailable, + FreeCpu: 4000, + FreeMemory: 8000, + Priority: 7, + } + bestFitWorker := &types.Worker{ + Id: "best-fit", + Status: types.WorkerStatusAvailable, + FreeCpu: 1500, + FreeMemory: 2000, + Priority: 7, + } + + worker, err := wb.selectWorkerFromWorkers([]*types.Worker{roomyWorker, bestFitWorker}, request) + assert.Nil(t, err) + assert.Equal(t, test.expectedID, worker.Id) + }) + } } -func TestSelectWorkerBestFitPreservesPriority(t *testing.T) { +func TestSelectWorkerBestFitPreservesPriorityWhenEnabled(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) + wb.config.Worker.PreferBestFit = true request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} highPriorityWorker := &types.Worker{ @@ -2902,9 +2973,10 @@ func TestSelectWorkerBestFitPreservesPriority(t *testing.T) { assert.Equal(t, highPriorityWorker.Id, worker.Id) } -func TestSelectWorkerBestFitPreservesAvailability(t *testing.T) { +func TestSelectWorkerBestFitPreservesAvailabilityWhenEnabled(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) + wb.config.Worker.PreferBestFit = true request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} availableWorker := &types.Worker{ @@ -2932,9 +3004,10 @@ func TestSelectWorkerBestFitPreservesAvailability(t *testing.T) { assert.Equal(t, availableWorker.Id, worker.Id) } -func TestSelectWorkerBestFitUsesWorkerIDForDeterministicTies(t *testing.T) { +func TestSelectWorkerBestFitUsesWorkerIDForDeterministicTiesWhenEnabled(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) + wb.config.Worker.PreferBestFit = true request := &types.ContainerRequest{Cpu: 1000, Memory: 1000} workerB := &types.Worker{ diff --git a/pkg/types/config.go b/pkg/types/config.go index 629cbe26d..4980e8f70 100644 --- a/pkg/types/config.go +++ b/pkg/types/config.go @@ -407,6 +407,7 @@ type WorkerConfig struct { Namespace string `key:"namespace" json:"namespace"` ServiceAccountName string `key:"serviceAccountName" json:"service_account_name"` JobResourcesEnforced bool `key:"jobResourcesEnforced" json:"job_resources_enforced"` + PreferBestFit bool `key:"preferBestFit" json:"prefer_best_fit"` ContainerResourceLimits ContainerResourceLimitsConfig `key:"containerResourceLimits" json:"container_resource_limits"` DefaultWorkerCPURequest int64 `key:"defaultWorkerCPURequest" json:"default_worker_cpu_request"` DefaultWorkerMemoryRequest int64 `key:"defaultWorkerMemoryRequest" json:"default_worker_memory_request"` From 19dff61ff5ce61ab0fda75ea62c7ca57ec29ad63 Mon Sep 17 00:00:00 2001 From: Luke Lombardi <33990301+luke-lombardi@users.noreply.github.com> Date: Wed, 29 Jul 2026 13:07:45 -0400 Subject: [PATCH 5/5] fix: keep best-fit packing disabled by default --- pkg/common/config.default.yaml | 1 + pkg/scheduler/scheduler_test.go | 17 +++++++++++++---- 2 files changed, 14 insertions(+), 4 deletions(-) diff --git a/pkg/common/config.default.yaml b/pkg/common/config.default.yaml index 28c3526a3..390a56423 100644 --- a/pkg/common/config.default.yaml +++ b/pkg/common/config.default.yaml @@ -285,6 +285,7 @@ worker: # non-standard k8s job spec imagePVCName: beta9-images jobResourcesEnforced: false + preferBestFit: false containerResourceLimits: cpuEnforced: true cpuAffinityEnforced: false diff --git a/pkg/scheduler/scheduler_test.go b/pkg/scheduler/scheduler_test.go index ad2534787..457eb8137 100644 --- a/pkg/scheduler/scheduler_test.go +++ b/pkg/scheduler/scheduler_test.go @@ -2900,12 +2900,18 @@ func TestSelectGPUWorkerDoesNotMutatePriority(t *testing.T) { func TestSelectWorkerCapacityTieBreakHonorsPackingPreference(t *testing.T) { tests := []struct { - name string - enableBestFit bool - expectedID string + name string + omitWorkerConfig bool + enableBestFit bool + expectedID string }{ { - name: "spread by default", + name: "spread when setting is omitted", + omitWorkerConfig: true, + expectedID: "roomy", + }, + { + name: "spread from default config", expectedID: "roomy", }, { @@ -2920,6 +2926,9 @@ func TestSelectWorkerCapacityTieBreakHonorsPackingPreference(t *testing.T) { wb, err := NewSchedulerForTest() assert.Nil(t, err) assert.False(t, wb.config.Worker.PreferBestFit) + if test.omitWorkerConfig { + wb.config.Worker = types.WorkerConfig{} + } if test.enableBestFit { wb.config.Worker.PreferBestFit = true }