From e3144dde807d4cf8172411880c692b0d08598ed6 Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Wed, 29 Jul 2026 00:54:57 +0200 Subject: [PATCH 01/13] Add resourceaware plugin: configurable resource-based JobOrderFn tiebreak Signed-off-by: CoolingCube --- .../plugins/resourceaware/resourceaware.go | 56 +++++++++++ .../resourceaware/resourceaware_test.go | 92 +++++++++++++++++++ 2 files changed, 148 insertions(+) create mode 100644 pkg/scheduler/plugins/resourceaware/resourceaware.go create mode 100644 pkg/scheduler/plugins/resourceaware/resourceaware_test.go diff --git a/pkg/scheduler/plugins/resourceaware/resourceaware.go b/pkg/scheduler/plugins/resourceaware/resourceaware.go new file mode 100644 index 000000000..446db857c --- /dev/null +++ b/pkg/scheduler/plugins/resourceaware/resourceaware.go @@ -0,0 +1,56 @@ +package resourceaware + +import ( + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/log" +) + +type resourceAwarePlugin struct { + mode string +} + +func New(arguments framework.PluginArguments) framework.Plugin { + mode := arguments.GetString("mode", "prefer-larger") + if mode != "prefer-larger" && mode != "prefer-smaller" { + log.InfraLogger.Warningf("resourceaware: unrecognized mode %q, defaulting to prefer-larger", mode) + mode = "prefer-larger" + } + return &resourceAwarePlugin{mode: mode} +} + +func (rp *resourceAwarePlugin) Name() string { + return "resourceaware" +} + +func (rp *resourceAwarePlugin) OnSessionOpen(ssn *framework.Session) { + ssn.AddJobOrderFn(rp.JobOrderFn) +} + +func (rp *resourceAwarePlugin) JobOrderFn(l, r interface{}) int { + lv := l.(*podgroup_info.PodGroupInfo) + rv := r.(*podgroup_info.PodGroupInfo) + + lGPU := lv.GetAliveTasksRequestedGPUs() + rGPU := rv.GetAliveTasksRequestedGPUs() + + switch rp.mode { + case "prefer-smaller": + if lGPU < rGPU { + return -1 + } + if lGPU > rGPU { + return 1 + } + default: // "prefer-larger" + if lGPU > rGPU { + return -1 + } + if lGPU < rGPU { + return 1 + } + } + return 0 +} + +func (rp *resourceAwarePlugin) OnSessionClose(_ *framework.Session) {} \ No newline at end of file diff --git a/pkg/scheduler/plugins/resourceaware/resourceaware_test.go b/pkg/scheduler/plugins/resourceaware/resourceaware_test.go new file mode 100644 index 000000000..972e72922 --- /dev/null +++ b/pkg/scheduler/plugins/resourceaware/resourceaware_test.go @@ -0,0 +1,92 @@ +package resourceaware + +import ( + "testing" + + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/common_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/pod_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/pod_status" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/resource_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" +) + +func makeGPUPodGroup(uid string, priority int32, gpuCount float64, vm *resource_info.ResourceVectorMap) *podgroup_info.PodGroupInfo { + task := &pod_info.PodInfo{ + UID: common_info.PodID(uid + "-task"), + ResReqVector: resource_info.NewResourceVectorWithValues(0, 0, gpuCount, vm), + Status: pod_status.Running, + } + pg := podgroup_info.NewPodGroupInfoWithVectorMap(common_info.PodGroupID(uid), vm, task) + pg.Priority = priority + return pg +} + +func newPlugin(t *testing.T, mode string) *resourceAwarePlugin { + t.Helper() + args := framework.PluginArguments{} + if mode != "" { + args["mode"] = mode + } + rp, ok := New(args).(*resourceAwarePlugin) + if !ok { + t.Fatalf("New() did not return *resourceAwarePlugin") + } + return rp +} + +func TestJobOrderFn_PriorityDiffers_Defers(t *testing.T) { + vm := resource_info.NewResourceVectorMap() + a := makeGPUPodGroup("a", 50, 1, vm) + b := makeGPUPodGroup("b", 10, 1, vm) + rp := newPlugin(t, "") + if got := rp.JobOrderFn(a, b); got != 0 { + t.Errorf("expected 0 when priorities differ, got %d", got) + } +} + +func TestJobOrderFn_SamePriority_PrefersLarger_Default(t *testing.T) { + vm := resource_info.NewResourceVectorMap() + small := makeGPUPodGroup("small", 10, 1, vm) + large := makeGPUPodGroup("large", 10, 2, vm) + rp := newPlugin(t, "") + if got := rp.JobOrderFn(large, small); got != -1 { + t.Errorf("expected -1 (larger job preferred as victim), got %d", got) + } + if got := rp.JobOrderFn(small, large); got != 1 { + t.Errorf("expected 1, got %d", got) + } +} + +func TestJobOrderFn_SamePriority_SameGPU_FallsThrough(t *testing.T) { + vm := resource_info.NewResourceVectorMap() + a := makeGPUPodGroup("a", 10, 1, vm) + b := makeGPUPodGroup("b", 10, 1, vm) + rp := newPlugin(t, "") + if got := rp.JobOrderFn(a, b); got != 0 { + t.Errorf("expected 0 (equal GPU falls through), got %d", got) + } +} + +func TestJobOrderFn_PreferSmallerMode(t *testing.T) { + vm := resource_info.NewResourceVectorMap() + small := makeGPUPodGroup("small", 10, 1, vm) + large := makeGPUPodGroup("large", 10, 2, vm) + rp := newPlugin(t, "prefer-smaller") + if got := rp.JobOrderFn(small, large); got != -1 { + t.Errorf("expected -1 (smaller job preferred as victim in prefer-smaller mode), got %d", got) + } + if got := rp.JobOrderFn(large, small); got != 1 { + t.Errorf("expected 1, got %d", got) + } +} + +func TestNew_UnrecognizedMode_FallsBackToPreferLarger(t *testing.T) { + vm := resource_info.NewResourceVectorMap() + small := makeGPUPodGroup("small", 10, 1, vm) + large := makeGPUPodGroup("large", 10, 2, vm) + rp := newPlugin(t, "totally-not-a-real-mode") + if got := rp.JobOrderFn(large, small); got != -1 { + t.Errorf("expected fallback to prefer-larger behavior, got %d", got) + } +} \ No newline at end of file From 60f765cbe1bb939ac2d12c71d897c6815f59c23c Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Wed, 29 Jul 2026 13:26:15 +0200 Subject: [PATCH 02/13] Register resourceaware plugin in factory.go; restore priority guard clause Signed-off-by: CoolingCube --- pkg/scheduler/plugins/factory.go | 4 +++- pkg/scheduler/plugins/resourceaware/resourceaware.go | 6 +++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/pkg/scheduler/plugins/factory.go b/pkg/scheduler/plugins/factory.go index 23b0a4486..6a69d6fb8 100644 --- a/pkg/scheduler/plugins/factory.go +++ b/pkg/scheduler/plugins/factory.go @@ -45,12 +45,14 @@ import ( "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/subgrouporder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/taskorder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/topology" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/resourceaware" ) func InitDefaultPlugins() { // Plugins for PodGroupInfos framework.RegisterPluginBuilder("predicates", predicates.New) framework.RegisterPluginBuilder("priority", priority.New) + framework.RegisterPluginBuilder("resourceaware", resourceaware.New) framework.RegisterPluginBuilder("nodeplacement", nodeplacement.New) framework.RegisterPluginBuilder("nominatednode", nominatednode.New) framework.RegisterPluginBuilder("numa", numa.New) @@ -79,4 +81,4 @@ func InitDefaultPlugins() { // Always register the Job Order Plugin last. framework.RegisterPluginBuilder("reflectjoborder", reflectjoborder.New) -} +} \ No newline at end of file diff --git a/pkg/scheduler/plugins/resourceaware/resourceaware.go b/pkg/scheduler/plugins/resourceaware/resourceaware.go index 446db857c..469c59fe2 100644 --- a/pkg/scheduler/plugins/resourceaware/resourceaware.go +++ b/pkg/scheduler/plugins/resourceaware/resourceaware.go @@ -31,6 +31,10 @@ func (rp *resourceAwarePlugin) JobOrderFn(l, r interface{}) int { lv := l.(*podgroup_info.PodGroupInfo) rv := r.(*podgroup_info.PodGroupInfo) + if lv.Priority != rv.Priority { + return 0 + } + lGPU := lv.GetAliveTasksRequestedGPUs() rGPU := rv.GetAliveTasksRequestedGPUs() @@ -53,4 +57,4 @@ func (rp *resourceAwarePlugin) JobOrderFn(l, r interface{}) int { return 0 } -func (rp *resourceAwarePlugin) OnSessionClose(_ *framework.Session) {} \ No newline at end of file +func (rp *resourceAwarePlugin) OnSessionClose(_ *framework.Session) {} From c5ff98c4128cb62d6708123d23a4297548524426 Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Wed, 29 Jul 2026 14:41:10 +0200 Subject: [PATCH 03/13] chore: add changelog fragment for resourceaware plugin Signed-off-by: CoolingCube --- .changes/unreleased/added-20260729-142722.yaml | 3 +++ 1 file changed, 3 insertions(+) create mode 100644 .changes/unreleased/added-20260729-142722.yaml diff --git a/.changes/unreleased/added-20260729-142722.yaml b/.changes/unreleased/added-20260729-142722.yaml new file mode 100644 index 000000000..56175d14a --- /dev/null +++ b/.changes/unreleased/added-20260729-142722.yaml @@ -0,0 +1,3 @@ +kind: Added +body: |- + Add resourceaware plugin for configurable GPU-based JobOrderFn tiebreak on priority ties From 279043fad4575a6720fd22af303b216865a0ee8e Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Thu, 30 Jul 2026 15:39:07 +0200 Subject: [PATCH 04/13] refactor(scheduler): rename resourceaware plugin to gpujoborder for clarity Signed-off-by: CoolingCube --- pkg/scheduler/plugins/factory.go | 4 +- .../plugins/gpujoborder/gpujoborder.go | 65 +++++++++++++++++++ .../gpujoborder_test.go} | 8 +-- .../plugins/resourceaware/resourceaware.go | 60 ----------------- 4 files changed, 71 insertions(+), 66 deletions(-) create mode 100644 pkg/scheduler/plugins/gpujoborder/gpujoborder.go rename pkg/scheduler/plugins/{resourceaware/resourceaware_test.go => gpujoborder/gpujoborder_test.go} (94%) delete mode 100644 pkg/scheduler/plugins/resourceaware/resourceaware.go diff --git a/pkg/scheduler/plugins/factory.go b/pkg/scheduler/plugins/factory.go index 6a69d6fb8..aa5b49b07 100644 --- a/pkg/scheduler/plugins/factory.go +++ b/pkg/scheduler/plugins/factory.go @@ -45,14 +45,14 @@ import ( "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/subgrouporder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/taskorder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/topology" - "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/resourceaware" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpujoborder" ) func InitDefaultPlugins() { // Plugins for PodGroupInfos framework.RegisterPluginBuilder("predicates", predicates.New) framework.RegisterPluginBuilder("priority", priority.New) - framework.RegisterPluginBuilder("resourceaware", resourceaware.New) + framework.RegisterPluginBuilder("gpujoborder", gpujoborder.New) framework.RegisterPluginBuilder("nodeplacement", nodeplacement.New) framework.RegisterPluginBuilder("nominatednode", nominatednode.New) framework.RegisterPluginBuilder("numa", numa.New) diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go new file mode 100644 index 000000000..1c667eb02 --- /dev/null +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go @@ -0,0 +1,65 @@ +package gpujoborder + +import ( + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/log" +) + +const ( + ModePreferLarger = "prefer-larger" + ModePreferSmaller = "prefer-smaller" +) + +type gpuJobOrderPlugin struct { + mode string +} + +func New(arguments framework.PluginArguments) framework.Plugin { + mode := arguments.GetString("mode", ModePreferLarger) + if mode != ModePreferLarger && mode != ModePreferSmaller { + log.InfraLogger.Warningf("gpujoborder: unrecognized mode %q, defaulting to prefer-larger", mode) + mode = ModePreferLarger + } + return &gpuJobOrderPlugin{mode: mode} +} + +func (rp *gpuJobOrderPlugin) Name() string { + return "gpujoborder" +} + +func (rp *gpuJobOrderPlugin) OnSessionOpen(ssn *framework.Session) { + ssn.AddJobOrderFn(rp.JobOrderFn) +} + +func (rp *gpuJobOrderPlugin) JobOrderFn(l, r interface{}) int { + lv := l.(*podgroup_info.PodGroupInfo) + rv := r.(*podgroup_info.PodGroupInfo) + + if lv.Priority != rv.Priority { + return 0 + } + + lGPU := lv.GetAliveTasksRequestedGPUs() + rGPU := rv.GetAliveTasksRequestedGPUs() + + switch rp.mode { + case ModePreferSmaller: + if lGPU < rGPU { + return -1 + } + if lGPU > rGPU { + return 1 + } + default: // ModePreferLarger + if lGPU > rGPU { + return -1 + } + if lGPU < rGPU { + return 1 + } + } + return 0 +} + +func (rp *gpuJobOrderPlugin) OnSessionClose(_ *framework.Session) {} \ No newline at end of file diff --git a/pkg/scheduler/plugins/resourceaware/resourceaware_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go similarity index 94% rename from pkg/scheduler/plugins/resourceaware/resourceaware_test.go rename to pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index 972e72922..a2bfaa6e1 100644 --- a/pkg/scheduler/plugins/resourceaware/resourceaware_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -1,4 +1,4 @@ -package resourceaware +package gpujoborder import ( "testing" @@ -22,15 +22,15 @@ func makeGPUPodGroup(uid string, priority int32, gpuCount float64, vm *resource_ return pg } -func newPlugin(t *testing.T, mode string) *resourceAwarePlugin { +func newPlugin(t *testing.T, mode string) *gpuJobOrderPlugin { t.Helper() args := framework.PluginArguments{} if mode != "" { args["mode"] = mode } - rp, ok := New(args).(*resourceAwarePlugin) + rp, ok := New(args).(*gpuJobOrderPlugin) if !ok { - t.Fatalf("New() did not return *resourceAwarePlugin") + t.Fatalf("New() did not return *gpuJobOrderPlugin") } return rp } diff --git a/pkg/scheduler/plugins/resourceaware/resourceaware.go b/pkg/scheduler/plugins/resourceaware/resourceaware.go deleted file mode 100644 index 469c59fe2..000000000 --- a/pkg/scheduler/plugins/resourceaware/resourceaware.go +++ /dev/null @@ -1,60 +0,0 @@ -package resourceaware - -import ( - "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" - "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" - "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/log" -) - -type resourceAwarePlugin struct { - mode string -} - -func New(arguments framework.PluginArguments) framework.Plugin { - mode := arguments.GetString("mode", "prefer-larger") - if mode != "prefer-larger" && mode != "prefer-smaller" { - log.InfraLogger.Warningf("resourceaware: unrecognized mode %q, defaulting to prefer-larger", mode) - mode = "prefer-larger" - } - return &resourceAwarePlugin{mode: mode} -} - -func (rp *resourceAwarePlugin) Name() string { - return "resourceaware" -} - -func (rp *resourceAwarePlugin) OnSessionOpen(ssn *framework.Session) { - ssn.AddJobOrderFn(rp.JobOrderFn) -} - -func (rp *resourceAwarePlugin) JobOrderFn(l, r interface{}) int { - lv := l.(*podgroup_info.PodGroupInfo) - rv := r.(*podgroup_info.PodGroupInfo) - - if lv.Priority != rv.Priority { - return 0 - } - - lGPU := lv.GetAliveTasksRequestedGPUs() - rGPU := rv.GetAliveTasksRequestedGPUs() - - switch rp.mode { - case "prefer-smaller": - if lGPU < rGPU { - return -1 - } - if lGPU > rGPU { - return 1 - } - default: // "prefer-larger" - if lGPU > rGPU { - return -1 - } - if lGPU < rGPU { - return 1 - } - } - return 0 -} - -func (rp *resourceAwarePlugin) OnSessionClose(_ *framework.Session) {} From d366fa3d4b6bb155ad3bc91ee14dc03e51bb520f Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Thu, 30 Jul 2026 17:25:02 +0200 Subject: [PATCH 05/13] docs: update changelog fragment to reference gpujoborder Signed-off-by: CoolingCube --- .changes/unreleased/added-20260729-142722.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changes/unreleased/added-20260729-142722.yaml b/.changes/unreleased/added-20260729-142722.yaml index 56175d14a..9c12dfcf1 100644 --- a/.changes/unreleased/added-20260729-142722.yaml +++ b/.changes/unreleased/added-20260729-142722.yaml @@ -1,3 +1,3 @@ kind: Added body: |- - Add resourceaware plugin for configurable GPU-based JobOrderFn tiebreak on priority ties + Add gpujoborder plugin for configurable GPU-based JobOrderFn tiebreak on priority ties From ff245eaf11416b634118ee12acead5a3f6fe68dd Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Thu, 30 Jul 2026 23:25:58 +0200 Subject: [PATCH 06/13] perf(scheduler): cache GetAliveTasksRequestedGPUs to fix reclaim benchmark regression Signed-off-by: CoolingCube --- pkg/scheduler/api/podgroup_info/job_info.go | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/pkg/scheduler/api/podgroup_info/job_info.go b/pkg/scheduler/api/podgroup_info/job_info.go index 1954839bf..a95b5707e 100644 --- a/pkg/scheduler/api/podgroup_info/job_info.go +++ b/pkg/scheduler/api/podgroup_info/job_info.go @@ -93,6 +93,7 @@ type PodGroupInfo struct { tasksToAllocateInitResourceVector resource_info.ResourceVector PodStatusIndex map[pod_status.PodStatus]pod_info.PodsMap activeAllocatedCount *int + aliveTasksRequestedGPUs *float64 } func NewPodGroupInfo(uid common_info.PodGroupID, tasks ...*pod_info.PodInfo) *PodGroupInfo { @@ -365,6 +366,7 @@ func (pgi *PodGroupInfo) invalidateTasksCache() { pgi.allPodsMap = nil pgi.tasksToAllocate = nil pgi.tasksToAllocateInitResourceVector = nil + pgi.aliveTasksRequestedGPUs = nil } func (pgi *PodGroupInfo) GetActiveAllocatedTasksCount() int { @@ -456,14 +458,16 @@ func (pgi *PodGroupInfo) GetNumGatedTasks() int { } func (pgi *PodGroupInfo) GetAliveTasksRequestedGPUs() float64 { - tasksTotalRequestedGPUs := float64(0) - for _, task := range pgi.GetAllPodsMap() { - if pod_status.IsAliveStatus(task.Status) { - tasksTotalRequestedGPUs += task.ResReqVector.Get(resource_info.GPUIndex) + if pgi.aliveTasksRequestedGPUs == nil { + tasksTotalRequestedGPUs := float64(0) + for _, task := range pgi.GetAllPodsMap() { + if pod_status.IsAliveStatus(task.Status) { + tasksTotalRequestedGPUs += task.ResReqVector.Get(resource_info.GPUIndex) + } } + pgi.aliveTasksRequestedGPUs = ptr.To(tasksTotalRequestedGPUs) } - - return tasksTotalRequestedGPUs + return *pgi.aliveTasksRequestedGPUs } func (pgi *PodGroupInfo) GetTasksActiveAllocatedReqResourceVector() resource_info.ResourceVector { @@ -672,4 +676,4 @@ func (pgi *PodGroupInfo) addInvalidSubGroupTask(ti *pod_info.PodInfo, taskSubGro pgi.Name, )) pgi.AddTaskFitErrors(ti, fitErrors) -} +} \ No newline at end of file From 8a53739dcd8710d70fa03df70c84798e28f9cb9f Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Fri, 31 Jul 2026 13:11:36 +0200 Subject: [PATCH 07/13] chore: add license headers and fix formatting to pass make validate Signed-off-by: CoolingCube --- pkg/scheduler/api/podgroup_info/job_info.go | 2 +- pkg/scheduler/plugins/factory.go | 4 ++-- pkg/scheduler/plugins/gpujoborder/gpujoborder.go | 5 ++++- pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go | 5 ++++- 4 files changed, 11 insertions(+), 5 deletions(-) diff --git a/pkg/scheduler/api/podgroup_info/job_info.go b/pkg/scheduler/api/podgroup_info/job_info.go index a95b5707e..9c468e13d 100644 --- a/pkg/scheduler/api/podgroup_info/job_info.go +++ b/pkg/scheduler/api/podgroup_info/job_info.go @@ -676,4 +676,4 @@ func (pgi *PodGroupInfo) addInvalidSubGroupTask(ti *pod_info.PodInfo, taskSubGro pgi.Name, )) pgi.AddTaskFitErrors(ti, fitErrors) -} \ No newline at end of file +} diff --git a/pkg/scheduler/plugins/factory.go b/pkg/scheduler/plugins/factory.go index aa5b49b07..6fcadcc55 100644 --- a/pkg/scheduler/plugins/factory.go +++ b/pkg/scheduler/plugins/factory.go @@ -23,6 +23,7 @@ import ( "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/dynamicresources" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/elastic" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpujoborder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpupack" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpusharingorder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpuspread" @@ -45,7 +46,6 @@ import ( "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/subgrouporder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/taskorder" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/topology" - "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/gpujoborder" ) func InitDefaultPlugins() { @@ -81,4 +81,4 @@ func InitDefaultPlugins() { // Always register the Job Order Plugin last. framework.RegisterPluginBuilder("reflectjoborder", reflectjoborder.New) -} \ No newline at end of file +} diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go index 1c667eb02..76e7751cc 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go @@ -1,3 +1,6 @@ +// Copyright 2026 NVIDIA CORPORATION +// SPDX-License-Identifier: Apache-2.0 + package gpujoborder import ( @@ -62,4 +65,4 @@ func (rp *gpuJobOrderPlugin) JobOrderFn(l, r interface{}) int { return 0 } -func (rp *gpuJobOrderPlugin) OnSessionClose(_ *framework.Session) {} \ No newline at end of file +func (rp *gpuJobOrderPlugin) OnSessionClose(_ *framework.Session) {} diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index a2bfaa6e1..9e525afc4 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -1,3 +1,6 @@ +// Copyright 2026 NVIDIA CORPORATION +// SPDX-License-Identifier: Apache-2.0 + package gpujoborder import ( @@ -89,4 +92,4 @@ func TestNew_UnrecognizedMode_FallsBackToPreferLarger(t *testing.T) { if got := rp.JobOrderFn(large, small); got != -1 { t.Errorf("expected fallback to prefer-larger behavior, got %d", got) } -} \ No newline at end of file +} From c7b3e9ef25eb8f390901ab73e298c2a2b2f422c5 Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Mon, 10 Aug 2026 18:34:44 +0200 Subject: [PATCH 08/13] fix(scheduler): scope gpujoborder to victim selection only via new VictimOrderFn extension point Fixes a sign-inversion bug where prefer-larger evicted the smaller job (the victim queue's !JobOrderFn inversion flipped the same sign used for pending-job ordering). Adds Session.VictimOrderFn/AddVictimOrderFn as a dedicated, non-inverted composition path for victim-specific comparators, falling back to the existing !JobOrderFn behavior when none are registered, so no other plugin's behavior changes. gpujoborder now registers exclusively via AddVictimOrderFn and no longer affects pending-job allocation ordering at all. Addresses @gshaibi's review comment on #1995. Signed-off-by: CoolingCube --- .../actions/utils/job_order_by_queue.go | 2 +- pkg/scheduler/framework/session.go | 2 + pkg/scheduler/framework/session_plugins.go | 24 ++++++++ .../plugins/gpujoborder/gpujoborder.go | 12 +++- .../plugins/gpujoborder/gpujoborder_test.go | 59 +++++++++++++++---- 5 files changed, 85 insertions(+), 14 deletions(-) diff --git a/pkg/scheduler/actions/utils/job_order_by_queue.go b/pkg/scheduler/actions/utils/job_order_by_queue.go index d8ca3c123..ce98df949 100644 --- a/pkg/scheduler/actions/utils/job_order_by_queue.go +++ b/pkg/scheduler/actions/utils/job_order_by_queue.go @@ -255,7 +255,7 @@ func (jo *JobsOrderByQueues) createLeafNode(queue *queue_info.QueueInfo) *queueN queue: queue, children: scheduler_util.NewPriorityQueue(func(l, r interface{}) bool { if jo.options.VictimQueue { - return !jo.ssn.JobOrderFn(l, r) + return jo.ssn.VictimOrderFn(l, r) } return jo.ssn.JobOrderFn(l, r) }, jo.options.MaxJobsQueueDepth), diff --git a/pkg/scheduler/framework/session.go b/pkg/scheduler/framework/session.go index 099a05548..37eb64def 100644 --- a/pkg/scheduler/framework/session.go +++ b/pkg/scheduler/framework/session.go @@ -76,6 +76,7 @@ type Session struct { NodePreOrderFns []api.NodePreOrderFn NodeOrderFns []api.NodeOrderFn JobOrderFns []common_info.CompareFn + VictimOrderFns []common_info.CompareFn SubGroupOrderFns []common_info.CompareFn TaskOrderFns []common_info.CompareFn QueueOrderFns []api.CompareQueueFn @@ -434,6 +435,7 @@ func (ssn *Session) clear() { ssn.NodePreOrderFns = nil ssn.NodeOrderFns = nil ssn.JobOrderFns = nil + ssn.VictimOrderFns = nil ssn.SubGroupOrderFns = nil ssn.TaskOrderFns = nil ssn.QueueOrderFns = nil diff --git a/pkg/scheduler/framework/session_plugins.go b/pkg/scheduler/framework/session_plugins.go index 63c639e90..a0b996f83 100644 --- a/pkg/scheduler/framework/session_plugins.go +++ b/pkg/scheduler/framework/session_plugins.go @@ -68,6 +68,15 @@ func (ssn *Session) AddJobOrderFn(jof common_info.CompareFn) { ssn.JobOrderFns = append(ssn.JobOrderFns, jof) } +// AddVictimOrderFn registers a comparator that applies ONLY when ordering +// candidates for eviction (the victim queue), not for regular pending-job +// allocation ordering. Unlike JobOrderFn, no external inversion is applied +// to this comparator's result: a negative return means "l is the BETTER +// victim (should be evicted first)", directly. +func (ssn *Session) AddVictimOrderFn(vof common_info.CompareFn) { + ssn.VictimOrderFns = append(ssn.VictimOrderFns, vof) +} + func (ssn *Session) AddTaskOrderFn(tof common_info.CompareFn) { ssn.TaskOrderFns = append(ssn.TaskOrderFns, tof) } @@ -285,6 +294,21 @@ func (ssn *Session) JobOrderFn(l, r interface{}) bool { } } +// VictimOrderFn composes registered victim-specific comparators directly -- +// a negative result means "l is the BETTER victim", with no external +// inversion needed or applied (unlike the old !JobOrderFn(l, r) pattern). +// If no victim-specific comparators are registered, this falls back to the +// existing !JobOrderFn(l, r) behavior, preserving exact current behavior +// for any plugin that never needed the ordering/victim distinction. +func (ssn *Session) VictimOrderFn(l, r interface{}) bool { + for _, vof := range ssn.VictimOrderFns { + if v := vof(l, r); v != 0 { + return v < 0 + } + } + return !ssn.JobOrderFn(l, r) +} + func (ssn *Session) TaskOrderFn(l, r interface{}) bool { for _, compareTasks := range ssn.TaskOrderFns { if comparison := compareTasks(l, r); comparison != 0 { diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go index 76e7751cc..ef9bc81dc 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go @@ -31,11 +31,19 @@ func (rp *gpuJobOrderPlugin) Name() string { return "gpujoborder" } +// OnSessionOpen registers this plugin's comparator via AddVictimOrderFn, +// NOT AddJobOrderFn. This plugin is scoped to victim/eviction selection +// only, per the real, confirmed distinction between the ordering and +// victim-selection paths (see VictimOrderFn in session_plugins.go) -- +// pending-job allocation ordering is intentionally left untouched. func (rp *gpuJobOrderPlugin) OnSessionOpen(ssn *framework.Session) { - ssn.AddJobOrderFn(rp.JobOrderFn) + ssn.AddVictimOrderFn(rp.VictimOrderFn) } -func (rp *gpuJobOrderPlugin) JobOrderFn(l, r interface{}) int { +// VictimOrderFn returns -1 when l is the BETTER victim (should be evicted +// first), matching VictimOrderFn's direct (non-inverted) contract -- no +// external sign flip is applied or needed here. +func (rp *gpuJobOrderPlugin) VictimOrderFn(l, r interface{}) int { lv := l.(*podgroup_info.PodGroupInfo) rv := r.(*podgroup_info.PodGroupInfo) diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index 9e525afc4..ac4852749 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -38,48 +38,48 @@ func newPlugin(t *testing.T, mode string) *gpuJobOrderPlugin { return rp } -func TestJobOrderFn_PriorityDiffers_Defers(t *testing.T) { +func TestVictimOrderFn_PriorityDiffers_Defers(t *testing.T) { vm := resource_info.NewResourceVectorMap() a := makeGPUPodGroup("a", 50, 1, vm) b := makeGPUPodGroup("b", 10, 1, vm) rp := newPlugin(t, "") - if got := rp.JobOrderFn(a, b); got != 0 { + if got := rp.VictimOrderFn(a, b); got != 0 { t.Errorf("expected 0 when priorities differ, got %d", got) } } -func TestJobOrderFn_SamePriority_PrefersLarger_Default(t *testing.T) { +func TestVictimOrderFn_SamePriority_PrefersLarger_Default(t *testing.T) { vm := resource_info.NewResourceVectorMap() small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) rp := newPlugin(t, "") - if got := rp.JobOrderFn(large, small); got != -1 { + if got := rp.VictimOrderFn(large, small); got != -1 { t.Errorf("expected -1 (larger job preferred as victim), got %d", got) } - if got := rp.JobOrderFn(small, large); got != 1 { + if got := rp.VictimOrderFn(small, large); got != 1 { t.Errorf("expected 1, got %d", got) } } -func TestJobOrderFn_SamePriority_SameGPU_FallsThrough(t *testing.T) { +func TestVictimOrderFn_SamePriority_SameGPU_FallsThrough(t *testing.T) { vm := resource_info.NewResourceVectorMap() a := makeGPUPodGroup("a", 10, 1, vm) b := makeGPUPodGroup("b", 10, 1, vm) rp := newPlugin(t, "") - if got := rp.JobOrderFn(a, b); got != 0 { + if got := rp.VictimOrderFn(a, b); got != 0 { t.Errorf("expected 0 (equal GPU falls through), got %d", got) } } -func TestJobOrderFn_PreferSmallerMode(t *testing.T) { +func TestVictimOrderFn_PreferSmallerMode(t *testing.T) { vm := resource_info.NewResourceVectorMap() small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) rp := newPlugin(t, "prefer-smaller") - if got := rp.JobOrderFn(small, large); got != -1 { + if got := rp.VictimOrderFn(small, large); got != -1 { t.Errorf("expected -1 (smaller job preferred as victim in prefer-smaller mode), got %d", got) } - if got := rp.JobOrderFn(large, small); got != 1 { + if got := rp.VictimOrderFn(large, small); got != 1 { t.Errorf("expected 1, got %d", got) } } @@ -89,7 +89,44 @@ func TestNew_UnrecognizedMode_FallsBackToPreferLarger(t *testing.T) { small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) rp := newPlugin(t, "totally-not-a-real-mode") - if got := rp.JobOrderFn(large, small); got != -1 { + if got := rp.VictimOrderFn(large, small); got != -1 { t.Errorf("expected fallback to prefer-larger behavior, got %d", got) } } + +// TestGpujoborder_DoesNotAffect_PendingJobOrdering is the specific test +// gshaibi asked for: confirms that registering gpujoborder's comparator +// has NO EFFECT on ssn.JobOrderFn (pending-job allocation ordering), +// since the plugin now registers exclusively via AddVictimOrderFn. +func TestGpujoborder_DoesNotAffect_PendingJobOrdering(t *testing.T) { + rp := newPlugin(t, "") + + vm := resource_info.NewResourceVectorMap() + small := makeGPUPodGroup("small", 10, 1, vm) + large := makeGPUPodGroup("large", 10, 2, vm) + + // Session WITHOUT gpujoborder registered: falls through to the + // existing CreationTimestamp/UID fallback for JobOrderFn. + ssnWithout := &framework.Session{} + baselineResult := ssnWithout.JobOrderFn(small, large) + + // Session WITH gpujoborder registered via OnSessionOpen (the real + // registration path), same two jobs, same JobOrderFn call. + ssnWith := &framework.Session{} + rp.OnSessionOpen(ssnWith) + withPluginResult := ssnWith.JobOrderFn(small, large) + + if len(ssnWith.JobOrderFns) != 0 { + t.Errorf("expected gpujoborder to register 0 JobOrderFns (it should only "+ + "register via AddVictimOrderFn now), got %d", len(ssnWith.JobOrderFns)) + } + if len(ssnWith.VictimOrderFns) != 1 { + t.Errorf("expected gpujoborder to register exactly 1 VictimOrderFn, got %d", + len(ssnWith.VictimOrderFns)) + } + if baselineResult != withPluginResult { + t.Errorf("pending-job ordering changed after registering gpujoborder: "+ + "without=%v, with=%v -- expected identical, since gpujoborder must not "+ + "affect JobOrderFn at all", baselineResult, withPluginResult) + } +} From 5991c0558d62b80c7f7789b220b526adc90db9fe Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Tue, 11 Aug 2026 01:18:53 +0200 Subject: [PATCH 09/13] fix(scheduler): VictimOrderFn correctly composes with existing JobOrderFns; rename modes for clarity Addresses @gshaibi's two review comments: 1. VictimOrderFn now checks the JobOrderFns chain first (priority.go, elastic.go, etc.) and only falls through to VictimOrderFns as a tiebreak when no registered JobOrderFn has an opinion. Previously, VictimOrderFns were checked unconditionally first, letting a raw resource-size comparator outrank elastic's at-min/above-min protection. New test confirms elastic correctly outranks gpujoborder now. 2. Renamed prefer-larger/prefer-smaller to evict-larger-first/ evict-smaller-first to remove ambiguity about whether the mode means 'prefer to keep' or 'prefer to evict' the larger job. Added package and constant documentation. Signed-off-by: CoolingCube --- pkg/scheduler/framework/session_plugins.go | 59 +++++++--- .../plugins/gpujoborder/gpujoborder.go | 37 ++++-- .../plugins/gpujoborder/gpujoborder_test.go | 111 ++++++++++++++---- 3 files changed, 158 insertions(+), 49 deletions(-) diff --git a/pkg/scheduler/framework/session_plugins.go b/pkg/scheduler/framework/session_plugins.go index a0b996f83..becd53e99 100644 --- a/pkg/scheduler/framework/session_plugins.go +++ b/pkg/scheduler/framework/session_plugins.go @@ -277,36 +277,61 @@ func (ssn *Session) QueueAllocatedResources(queue *queue_info.QueueInfo) *resour return nil } -func (ssn *Session) JobOrderFn(l, r interface{}) bool { - for _, jof := range ssn.JobOrderFns { - if j := jof(l, r); j != 0 { - return j < 0 - } - } - - // If no job order funcs, order job by CreationTimestamp first, then by UID. +// jobOrderCreationFallback is the shared tiebreak used when no registered +// comparator has an opinion: CreationTimestamp first, then UID. Extracted +// so both JobOrderFn and VictimOrderFn can share it without duplicating +// the comparison logic. +func jobOrderCreationFallback(l, r interface{}) bool { lv := l.(*podgroup_info.PodGroupInfo) rv := r.(*podgroup_info.PodGroupInfo) if lv.CreationTimestamp.Equal(&rv.CreationTimestamp) { return lv.UID < rv.UID - } else { - return lv.CreationTimestamp.Before(&rv.CreationTimestamp) } + return lv.CreationTimestamp.Before(&rv.CreationTimestamp) } -// VictimOrderFn composes registered victim-specific comparators directly -- -// a negative result means "l is the BETTER victim", with no external -// inversion needed or applied (unlike the old !JobOrderFn(l, r) pattern). -// If no victim-specific comparators are registered, this falls back to the -// existing !JobOrderFn(l, r) behavior, preserving exact current behavior -// for any plugin that never needed the ordering/victim distinction. +func (ssn *Session) JobOrderFn(l, r interface{}) bool { + for _, jof := range ssn.JobOrderFns { + if j := jof(l, r); j != 0 { + return j < 0 + } + } + return jobOrderCreationFallback(l, r) +} + +// VictimOrderFn composes registered victim-specific comparators, but only +// as a tiebreak AFTER the existing JobOrderFns chain (priority.go, +// elastic.go, etc.) has had its say -- fixed per @gshaibi's review: the +// original version checked VictimOrderFns first unconditionally, which +// let a raw resource-size comparator (e.g. gpujoborder) outrank +// elastic's deliberate at-min/above-min protection. Since JobOrderFn +// itself always resolves via its own CreationTimestamp/UID fallback, the +// raw JobOrderFns slice is iterated directly here (not the composed +// method) so a genuine "no opinion" state (all registered JobOrderFns +// return 0) can be detected before falling through to victim-specific +// comparators. +// +// A negative result from a JobOrderFn means "l ordered first for +// allocation" -- inverted here (j > 0) since a job LESS preferred for +// allocation should be MORE preferred as a victim. VictimOrderFns use +// direct (non-inverted) semantics: negative means "l is the better +// victim". If neither chain has an opinion, falls back to the inverted +// creation-timestamp order, preserving old behavior exactly. +// +// Path without any VictimOrderFns registered is byte-identical to the +// original !JobOrderFn(l, r) behavior. func (ssn *Session) VictimOrderFn(l, r interface{}) bool { + for _, jof := range ssn.JobOrderFns { + if j := jof(l, r); j != 0 { + return j > 0 + } + } for _, vof := range ssn.VictimOrderFns { if v := vof(l, r); v != 0 { return v < 0 } } - return !ssn.JobOrderFn(l, r) + return !jobOrderCreationFallback(l, r) } func (ssn *Session) TaskOrderFn(l, r interface{}) bool { diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go index ef9bc81dc..f16bea07f 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go @@ -1,6 +1,13 @@ // Copyright 2026 NVIDIA CORPORATION // SPDX-License-Identifier: Apache-2.0 +// Package gpujoborder implements a GPU-count-based victim-selection +// tiebreak, used only when two jobs of equal priority are otherwise tied +// (including by any other registered JobOrderFn, such as elastic's +// at-min/above-min protection) and the scheduler must choose which one +// to evict. It has no effect on pending-job allocation ordering, and no +// effect whenever an existing JobOrderFn-registered plugin already has an +// opinion on the comparison. package gpujoborder import ( @@ -10,8 +17,16 @@ import ( ) const ( - ModePreferLarger = "prefer-larger" - ModePreferSmaller = "prefer-smaller" + // ModeEvictLargerFirst prefers evicting the job requesting MORE GPUs + // when a tie must be broken. Named for the eviction outcome directly + // (not "prefer-larger", which reads ambiguously as either "prefer to + // evict the larger job" or "prefer to keep/favor the larger job") -- + // per review discussion, this is deliberately unambiguous since the + // config string is hard to rename once real users depend on it. + ModeEvictLargerFirst = "evict-larger-first" + // ModeEvictSmallerFirst prefers evicting the job requesting FEWER + // GPUs when a tie must be broken. + ModeEvictSmallerFirst = "evict-smaller-first" ) type gpuJobOrderPlugin struct { @@ -19,10 +34,10 @@ type gpuJobOrderPlugin struct { } func New(arguments framework.PluginArguments) framework.Plugin { - mode := arguments.GetString("mode", ModePreferLarger) - if mode != ModePreferLarger && mode != ModePreferSmaller { - log.InfraLogger.Warningf("gpujoborder: unrecognized mode %q, defaulting to prefer-larger", mode) - mode = ModePreferLarger + mode := arguments.GetString("mode", ModeEvictLargerFirst) + if mode != ModeEvictLargerFirst && mode != ModeEvictSmallerFirst { + log.InfraLogger.Warningf("gpujoborder: unrecognized mode %q, defaulting to %s", mode, ModeEvictLargerFirst) + mode = ModeEvictLargerFirst } return &gpuJobOrderPlugin{mode: mode} } @@ -33,9 +48,9 @@ func (rp *gpuJobOrderPlugin) Name() string { // OnSessionOpen registers this plugin's comparator via AddVictimOrderFn, // NOT AddJobOrderFn. This plugin is scoped to victim/eviction selection -// only, per the real, confirmed distinction between the ordering and -// victim-selection paths (see VictimOrderFn in session_plugins.go) -- -// pending-job allocation ordering is intentionally left untouched. +// only, and only applies as a tiebreak after every registered JobOrderFn +// (priority.go, elastic.go, etc.) has already had a chance to decide -- +// see Session.VictimOrderFn for the real composition order. func (rp *gpuJobOrderPlugin) OnSessionOpen(ssn *framework.Session) { ssn.AddVictimOrderFn(rp.VictimOrderFn) } @@ -55,14 +70,14 @@ func (rp *gpuJobOrderPlugin) VictimOrderFn(l, r interface{}) int { rGPU := rv.GetAliveTasksRequestedGPUs() switch rp.mode { - case ModePreferSmaller: + case ModeEvictSmallerFirst: if lGPU < rGPU { return -1 } if lGPU > rGPU { return 1 } - default: // ModePreferLarger + default: // ModeEvictLargerFirst if lGPU > rGPU { return -1 } diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index ac4852749..c9507a895 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -10,8 +10,11 @@ import ( "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/pod_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/pod_status" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info/subgroup_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/resource_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/framework" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/plugins/elastic" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/scheduler_util" ) func makeGPUPodGroup(uid string, priority int32, gpuCount float64, vm *resource_info.ResourceVectorMap) *podgroup_info.PodGroupInfo { @@ -48,7 +51,7 @@ func TestVictimOrderFn_PriorityDiffers_Defers(t *testing.T) { } } -func TestVictimOrderFn_SamePriority_PrefersLarger_Default(t *testing.T) { +func TestVictimOrderFn_SamePriority_EvictLargerFirst_Default(t *testing.T) { vm := resource_info.NewResourceVectorMap() small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) @@ -71,33 +74,34 @@ func TestVictimOrderFn_SamePriority_SameGPU_FallsThrough(t *testing.T) { } } -func TestVictimOrderFn_PreferSmallerMode(t *testing.T) { +func TestVictimOrderFn_EvictSmallerFirstMode(t *testing.T) { vm := resource_info.NewResourceVectorMap() small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) - rp := newPlugin(t, "prefer-smaller") + rp := newPlugin(t, "evict-smaller-first") if got := rp.VictimOrderFn(small, large); got != -1 { - t.Errorf("expected -1 (smaller job preferred as victim in prefer-smaller mode), got %d", got) + t.Errorf("expected -1 (smaller job preferred as victim in evict-smaller-first mode), got %d", got) } if got := rp.VictimOrderFn(large, small); got != 1 { t.Errorf("expected 1, got %d", got) } } -func TestNew_UnrecognizedMode_FallsBackToPreferLarger(t *testing.T) { +func TestNew_UnrecognizedMode_FallsBackToEvictLargerFirst(t *testing.T) { vm := resource_info.NewResourceVectorMap() small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) rp := newPlugin(t, "totally-not-a-real-mode") if got := rp.VictimOrderFn(large, small); got != -1 { - t.Errorf("expected fallback to prefer-larger behavior, got %d", got) + t.Errorf("expected fallback to evict-larger-first behavior, got %d", got) } } -// TestGpujoborder_DoesNotAffect_PendingJobOrdering is the specific test -// gshaibi asked for: confirms that registering gpujoborder's comparator -// has NO EFFECT on ssn.JobOrderFn (pending-job allocation ordering), -// since the plugin now registers exclusively via AddVictimOrderFn. +// TestGpujoborder_DoesNotAffect_PendingJobOrdering confirms gpujoborder +// registers ONLY via AddVictimOrderFn, never AddJobOrderFn -- verifying +// pending-job allocation ordering is byte-identical with or without the +// plugin registered (the test gshaibi asked for on the original scoping +// review comment). func TestGpujoborder_DoesNotAffect_PendingJobOrdering(t *testing.T) { rp := newPlugin(t, "") @@ -105,28 +109,93 @@ func TestGpujoborder_DoesNotAffect_PendingJobOrdering(t *testing.T) { small := makeGPUPodGroup("small", 10, 1, vm) large := makeGPUPodGroup("large", 10, 2, vm) - // Session WITHOUT gpujoborder registered: falls through to the - // existing CreationTimestamp/UID fallback for JobOrderFn. ssnWithout := &framework.Session{} baselineResult := ssnWithout.JobOrderFn(small, large) - // Session WITH gpujoborder registered via OnSessionOpen (the real - // registration path), same two jobs, same JobOrderFn call. ssnWith := &framework.Session{} rp.OnSessionOpen(ssnWith) withPluginResult := ssnWith.JobOrderFn(small, large) if len(ssnWith.JobOrderFns) != 0 { - t.Errorf("expected gpujoborder to register 0 JobOrderFns (it should only "+ - "register via AddVictimOrderFn now), got %d", len(ssnWith.JobOrderFns)) + t.Errorf("expected gpujoborder to register 0 JobOrderFns, got %d", len(ssnWith.JobOrderFns)) } if len(ssnWith.VictimOrderFns) != 1 { - t.Errorf("expected gpujoborder to register exactly 1 VictimOrderFn, got %d", - len(ssnWith.VictimOrderFns)) + t.Errorf("expected gpujoborder to register exactly 1 VictimOrderFn, got %d", len(ssnWith.VictimOrderFns)) } if baselineResult != withPluginResult { - t.Errorf("pending-job ordering changed after registering gpujoborder: "+ - "without=%v, with=%v -- expected identical, since gpujoborder must not "+ - "affect JobOrderFn at all", baselineResult, withPluginResult) + t.Errorf("pending-job ordering changed after registering gpujoborder: without=%v, with=%v", + baselineResult, withPluginResult) + } +} + +// makeElasticJob builds a same-priority job with one subgroup, letting the +// caller specify how many tasks are allocated relative to the subgroup's +// MinAvailable -- used to construct genuine at-min / above-min states for +// elastic.JobOrderFn to evaluate. Uses the real constructor first (so +// activeAllocatedCount and other internal fields are properly +// initialized, avoiding a nil-pointer panic a raw struct literal would +// risk), then swaps in a custom-MinAvailable subgroup before adding tasks. +func makeElasticJob(uid string, priority int32, minAvailable int32, allocatedTasks int, + gpuPerTask float64, vm *resource_info.ResourceVectorMap) *podgroup_info.PodGroupInfo { + + pg := podgroup_info.NewPodGroupInfoWithVectorMap(common_info.PodGroupID(uid), vm) + pg.Priority = priority + + root := subgroup_info.NewSubGroupSet(subgroup_info.RootSubGroupSetName, nil) + root.AddPodSet(subgroup_info.NewPodSet(podgroup_info.DefaultSubGroup, minAvailable, nil)) + pg.RootSubGroupSet = root + pg.PodSets = root.GetDescendantPodSets() + + for i := 0; i < allocatedTasks; i++ { + task := &pod_info.PodInfo{ + UID: common_info.PodID(uid + "-task"), + ResReqVector: resource_info.NewResourceVectorWithValues(0, 0, gpuPerTask, vm), + Status: pod_status.Running, + } + pg.AddTaskInfo(task) + } + return pg +} + +// TestElasticProtection_OutranksGpujoborder is the test @gshaibi asked +// for: confirms that when both elastic.JobOrderFn and gpujoborder's +// VictimOrderFn are registered on a real Session, elastic's deliberate +// at-min/above-min protection correctly outranks gpujoborder's raw +// GPU-size comparison -- a smaller ABOVE-min job must be evicted before a +// larger AT-min job, even though gpujoborder alone would prefer the +// larger job as victim. +func TestElasticProtection_OutranksGpujoborder(t *testing.T) { + rp := newPlugin(t, "") // evict-larger-first (default) + + vm := resource_info.NewResourceVectorMap() + // At its minimum (2/2 allocated), 2 GPUs total -- LARGER, but must be + // protected per elastic's logic. + atMinLarge := makeElasticJob("at-min-large", 10, 2, 2, 1.0, vm) + // Above its minimum (2 allocated, min 1), 2 GPUs total via 2 tasks -- + // wait: to keep this a clean, unconfounded test, above-min job uses + // FEWER total GPUs than the at-min job, so gpujoborder ALONE would + // prefer evicting the at-min job (it's larger) -- the opposite of + // what elastic's protection requires. + aboveMinSmall := makeElasticJob("above-min-small", 10, 1, 2, 0.5, vm) + + ssn := &framework.Session{} + ssn.AddJobOrderFn(elastic.JobOrderFn) + ssn.AddVictimOrderFn(rp.VictimOrderFn) + + victimLessFn := func(l, r interface{}) bool { + return ssn.VictimOrderFn(l, r) + } + pq := scheduler_util.NewPriorityQueue(victimLessFn, scheduler_util.QueueCapacityInfinite) + pq.Push(atMinLarge) + pq.Push(aboveMinSmall) + + popped := pq.Pop().(*podgroup_info.PodGroupInfo) + + t.Logf("Popped as victim: %q (GPUs=%v)", popped.UID, popped.GetAliveTasksRequestedGPUs()) + + if popped.UID != aboveMinSmall.UID { + t.Errorf("expected elastic's above-min job to be evicted first (protecting the at-min job), "+ + "but got %q evicted instead -- gpujoborder's raw GPU-size comparison incorrectly "+ + "outranked elastic's protection", popped.UID) } } From 04467da1cf37d82e906464fa8712a3dc7e3e1dbe Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Tue, 11 Aug 2026 01:29:02 +0200 Subject: [PATCH 10/13] test(scheduler): fix duplicate task UIDs in elastic-composition test helper Each task in a multi-task subgroup was getting an identical UID, causing the second task to silently overwrite the first in the pod map used by GetAliveTasksRequestedGPUs(). Didn't affect the test's pass/fail correctness (elastic's own task-count tracking uses a separate counter, unaffected by the collision, and short-circuits before gpujoborder's comparator is ever consulted) but did produce a misleading debug log value. Fixed by giving each task a unique UID. Signed-off-by: CoolingCube --- pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index c9507a895..55b83950e 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -4,6 +4,7 @@ package gpujoborder import ( + "fmt" "testing" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/common_info" @@ -148,7 +149,7 @@ func makeElasticJob(uid string, priority int32, minAvailable int32, allocatedTas for i := 0; i < allocatedTasks; i++ { task := &pod_info.PodInfo{ - UID: common_info.PodID(uid + "-task"), + UID: common_info.PodID(fmt.Sprintf("%s-task-%d", uid, i)), ResReqVector: resource_info.NewResourceVectorWithValues(0, 0, gpuPerTask, vm), Status: pod_status.Running, } From 2b6eadecddd4bfb57abc541ce9cf782bcca652f5 Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Wed, 12 Aug 2026 14:35:43 +0200 Subject: [PATCH 11/13] style(scheduler): remove comments per @gshaibi's review preference Signed-off-by: CoolingCube --- pkg/scheduler/framework/session_plugins.go | 30 ------------------- .../plugins/gpujoborder/gpujoborder.go | 27 ++--------------- .../plugins/gpujoborder/gpujoborder_test.go | 28 +---------------- 3 files changed, 3 insertions(+), 82 deletions(-) diff --git a/pkg/scheduler/framework/session_plugins.go b/pkg/scheduler/framework/session_plugins.go index becd53e99..52494358f 100644 --- a/pkg/scheduler/framework/session_plugins.go +++ b/pkg/scheduler/framework/session_plugins.go @@ -68,11 +68,6 @@ func (ssn *Session) AddJobOrderFn(jof common_info.CompareFn) { ssn.JobOrderFns = append(ssn.JobOrderFns, jof) } -// AddVictimOrderFn registers a comparator that applies ONLY when ordering -// candidates for eviction (the victim queue), not for regular pending-job -// allocation ordering. Unlike JobOrderFn, no external inversion is applied -// to this comparator's result: a negative return means "l is the BETTER -// victim (should be evicted first)", directly. func (ssn *Session) AddVictimOrderFn(vof common_info.CompareFn) { ssn.VictimOrderFns = append(ssn.VictimOrderFns, vof) } @@ -277,10 +272,6 @@ func (ssn *Session) QueueAllocatedResources(queue *queue_info.QueueInfo) *resour return nil } -// jobOrderCreationFallback is the shared tiebreak used when no registered -// comparator has an opinion: CreationTimestamp first, then UID. Extracted -// so both JobOrderFn and VictimOrderFn can share it without duplicating -// the comparison logic. func jobOrderCreationFallback(l, r interface{}) bool { lv := l.(*podgroup_info.PodGroupInfo) rv := r.(*podgroup_info.PodGroupInfo) @@ -299,27 +290,6 @@ func (ssn *Session) JobOrderFn(l, r interface{}) bool { return jobOrderCreationFallback(l, r) } -// VictimOrderFn composes registered victim-specific comparators, but only -// as a tiebreak AFTER the existing JobOrderFns chain (priority.go, -// elastic.go, etc.) has had its say -- fixed per @gshaibi's review: the -// original version checked VictimOrderFns first unconditionally, which -// let a raw resource-size comparator (e.g. gpujoborder) outrank -// elastic's deliberate at-min/above-min protection. Since JobOrderFn -// itself always resolves via its own CreationTimestamp/UID fallback, the -// raw JobOrderFns slice is iterated directly here (not the composed -// method) so a genuine "no opinion" state (all registered JobOrderFns -// return 0) can be detected before falling through to victim-specific -// comparators. -// -// A negative result from a JobOrderFn means "l ordered first for -// allocation" -- inverted here (j > 0) since a job LESS preferred for -// allocation should be MORE preferred as a victim. VictimOrderFns use -// direct (non-inverted) semantics: negative means "l is the better -// victim". If neither chain has an opinion, falls back to the inverted -// creation-timestamp order, preserving old behavior exactly. -// -// Path without any VictimOrderFns registered is byte-identical to the -// original !JobOrderFn(l, r) behavior. func (ssn *Session) VictimOrderFn(l, r interface{}) bool { for _, jof := range ssn.JobOrderFns { if j := jof(l, r); j != 0 { diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go index f16bea07f..3f536baed 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder.go @@ -1,13 +1,6 @@ // Copyright 2026 NVIDIA CORPORATION // SPDX-License-Identifier: Apache-2.0 -// Package gpujoborder implements a GPU-count-based victim-selection -// tiebreak, used only when two jobs of equal priority are otherwise tied -// (including by any other registered JobOrderFn, such as elastic's -// at-min/above-min protection) and the scheduler must choose which one -// to evict. It has no effect on pending-job allocation ordering, and no -// effect whenever an existing JobOrderFn-registered plugin already has an -// opinion on the comparison. package gpujoborder import ( @@ -17,15 +10,7 @@ import ( ) const ( - // ModeEvictLargerFirst prefers evicting the job requesting MORE GPUs - // when a tie must be broken. Named for the eviction outcome directly - // (not "prefer-larger", which reads ambiguously as either "prefer to - // evict the larger job" or "prefer to keep/favor the larger job") -- - // per review discussion, this is deliberately unambiguous since the - // config string is hard to rename once real users depend on it. - ModeEvictLargerFirst = "evict-larger-first" - // ModeEvictSmallerFirst prefers evicting the job requesting FEWER - // GPUs when a tie must be broken. + ModeEvictLargerFirst = "evict-larger-first" ModeEvictSmallerFirst = "evict-smaller-first" ) @@ -46,18 +31,10 @@ func (rp *gpuJobOrderPlugin) Name() string { return "gpujoborder" } -// OnSessionOpen registers this plugin's comparator via AddVictimOrderFn, -// NOT AddJobOrderFn. This plugin is scoped to victim/eviction selection -// only, and only applies as a tiebreak after every registered JobOrderFn -// (priority.go, elastic.go, etc.) has already had a chance to decide -- -// see Session.VictimOrderFn for the real composition order. func (rp *gpuJobOrderPlugin) OnSessionOpen(ssn *framework.Session) { ssn.AddVictimOrderFn(rp.VictimOrderFn) } -// VictimOrderFn returns -1 when l is the BETTER victim (should be evicted -// first), matching VictimOrderFn's direct (non-inverted) contract -- no -// external sign flip is applied or needed here. func (rp *gpuJobOrderPlugin) VictimOrderFn(l, r interface{}) int { lv := l.(*podgroup_info.PodGroupInfo) rv := r.(*podgroup_info.PodGroupInfo) @@ -77,7 +54,7 @@ func (rp *gpuJobOrderPlugin) VictimOrderFn(l, r interface{}) int { if lGPU > rGPU { return 1 } - default: // ModeEvictLargerFirst + default: if lGPU > rGPU { return -1 } diff --git a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go index 55b83950e..bcccd9f27 100644 --- a/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go +++ b/pkg/scheduler/plugins/gpujoborder/gpujoborder_test.go @@ -98,11 +98,6 @@ func TestNew_UnrecognizedMode_FallsBackToEvictLargerFirst(t *testing.T) { } } -// TestGpujoborder_DoesNotAffect_PendingJobOrdering confirms gpujoborder -// registers ONLY via AddVictimOrderFn, never AddJobOrderFn -- verifying -// pending-job allocation ordering is byte-identical with or without the -// plugin registered (the test gshaibi asked for on the original scoping -// review comment). func TestGpujoborder_DoesNotAffect_PendingJobOrdering(t *testing.T) { rp := newPlugin(t, "") @@ -129,13 +124,6 @@ func TestGpujoborder_DoesNotAffect_PendingJobOrdering(t *testing.T) { } } -// makeElasticJob builds a same-priority job with one subgroup, letting the -// caller specify how many tasks are allocated relative to the subgroup's -// MinAvailable -- used to construct genuine at-min / above-min states for -// elastic.JobOrderFn to evaluate. Uses the real constructor first (so -// activeAllocatedCount and other internal fields are properly -// initialized, avoiding a nil-pointer panic a raw struct literal would -// risk), then swaps in a custom-MinAvailable subgroup before adding tasks. func makeElasticJob(uid string, priority int32, minAvailable int32, allocatedTasks int, gpuPerTask float64, vm *resource_info.ResourceVectorMap) *podgroup_info.PodGroupInfo { @@ -158,25 +146,11 @@ func makeElasticJob(uid string, priority int32, minAvailable int32, allocatedTas return pg } -// TestElasticProtection_OutranksGpujoborder is the test @gshaibi asked -// for: confirms that when both elastic.JobOrderFn and gpujoborder's -// VictimOrderFn are registered on a real Session, elastic's deliberate -// at-min/above-min protection correctly outranks gpujoborder's raw -// GPU-size comparison -- a smaller ABOVE-min job must be evicted before a -// larger AT-min job, even though gpujoborder alone would prefer the -// larger job as victim. func TestElasticProtection_OutranksGpujoborder(t *testing.T) { - rp := newPlugin(t, "") // evict-larger-first (default) + rp := newPlugin(t, "") vm := resource_info.NewResourceVectorMap() - // At its minimum (2/2 allocated), 2 GPUs total -- LARGER, but must be - // protected per elastic's logic. atMinLarge := makeElasticJob("at-min-large", 10, 2, 2, 1.0, vm) - // Above its minimum (2 allocated, min 1), 2 GPUs total via 2 tasks -- - // wait: to keep this a clean, unconfounded test, above-min job uses - // FEWER total GPUs than the at-min job, so gpujoborder ALONE would - // prefer evicting the at-min job (it's larger) -- the opposite of - // what elastic's protection requires. aboveMinSmall := makeElasticJob("above-min-small", 10, 1, 2, 0.5, vm) ssn := &framework.Session{} From 24d25872ef57b6ee307997f577c49c9037a3b41a Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Sun, 16 Aug 2026 20:30:20 +0200 Subject: [PATCH 12/13] style(scheduler): remove comment from new test, matching bare-minimum convention Signed-off-by: CoolingCube --- .../framework/session_plugins_test.go | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/pkg/scheduler/framework/session_plugins_test.go b/pkg/scheduler/framework/session_plugins_test.go index 73f3c598f..3d735d52f 100644 --- a/pkg/scheduler/framework/session_plugins_test.go +++ b/pkg/scheduler/framework/session_plugins_test.go @@ -18,10 +18,12 @@ import ( kaiv1 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1" "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/common_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/node_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/pod_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/podgroup_info/subgroup_info" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/resource_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/conf" ) @@ -305,3 +307,43 @@ func scenarioSearchDurationForTest(value string) metav1.Duration { } return metav1.Duration{Duration: duration} } + +func TestVictimOrderFn_PreservesExistingJobOrderFnBehavior(t *testing.T) { + priorityLikeJobOrderFn := func(l, r interface{}) int { + lv := l.(*podgroup_info.PodGroupInfo) + rv := r.(*podgroup_info.PodGroupInfo) + if lv.Priority > rv.Priority { + return -1 + } + if lv.Priority < rv.Priority { + return 1 + } + return 0 + } + + vm := resource_info.NewResourceVectorMap() + priorities := []int32{1, 1, 5, 5, 10, 20, 20, 100} + + for i, lp := range priorities { + for j, rp := range priorities { + if i == j { + continue + } + ssn := &Session{JobOrderFns: []common_info.CompareFn{priorityLikeJobOrderFn}} + lJob := podgroup_info.NewPodGroupInfoWithVectorMap( + common_info.PodGroupID("l"), vm) + lJob.Priority = lp + rJob := podgroup_info.NewPodGroupInfoWithVectorMap( + common_info.PodGroupID("r"), vm) + rJob.Priority = rp + + oldBehavior := !ssn.JobOrderFn(lJob, rJob) + newBehavior := ssn.VictimOrderFn(lJob, rJob) + + if oldBehavior != newBehavior { + t.Errorf("mismatch for priorities l=%d r=%d: old !JobOrderFn=%v, new VictimOrderFn=%v", + lp, rp, oldBehavior, newBehavior) + } + } + } +} From 99cd1c2f445139de808bc315bda4c7077f51741b Mon Sep 17 00:00:00 2001 From: CoolingCube Date: Mon, 17 Aug 2026 01:25:55 +0200 Subject: [PATCH 13/13] docs(scheduler): add doc comments for the VictimOrderFn extension point Signed-off-by: CoolingCube --- pkg/scheduler/framework/session_plugins.go | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/pkg/scheduler/framework/session_plugins.go b/pkg/scheduler/framework/session_plugins.go index 52494358f..ba3e1e930 100644 --- a/pkg/scheduler/framework/session_plugins.go +++ b/pkg/scheduler/framework/session_plugins.go @@ -68,6 +68,18 @@ func (ssn *Session) AddJobOrderFn(jof common_info.CompareFn) { ssn.JobOrderFns = append(ssn.JobOrderFns, jof) } +// AddVictimOrderFn registers a comparator used only to rank candidate +// eviction victims, without affecting pending-job allocation order. +// +// The registered JobOrderFns chain is always checked first (inverted, +// matching existing allocation-order semantics); comparators registered +// here apply only as a tiebreak, when every JobOrderFn treats l and r as +// equal. This ensures job-level protections (e.g. elastic at-min/above-min +// status) take precedence over any victim-specific comparator. +// +// Sign convention: vof(l, r) < 0 means l should be evicted before r. This +// is the direct, non-inverted sense -- unlike the JobOrderFn chain, which +// VictimOrderFn inverts internally when using it for the eviction path. func (ssn *Session) AddVictimOrderFn(vof common_info.CompareFn) { ssn.VictimOrderFns = append(ssn.VictimOrderFns, vof) } @@ -290,6 +302,12 @@ func (ssn *Session) JobOrderFn(l, r interface{}) bool { return jobOrderCreationFallback(l, r) } +// VictimOrderFn reports whether l should be evicted before r. +// +// The JobOrderFns chain is consulted first, inverted; VictimOrderFns +// registered via AddVictimOrderFn apply only when every JobOrderFn treats +// l and r as equal. With no VictimOrderFns registered, this reduces to +// the original !JobOrderFn(l, r) behavior. func (ssn *Session) VictimOrderFn(l, r interface{}) bool { for _, jof := range ssn.JobOrderFns { if j := jof(l, r); j != 0 {