diff --git a/.changes/unreleased/added-20260810-112224.yaml b/.changes/unreleased/added-20260810-112224.yaml new file mode 100644 index 000000000..9feeb1be6 --- /dev/null +++ b/.changes/unreleased/added-20260810-112224.yaml @@ -0,0 +1,3 @@ +kind: Added +body: |- + Enforce fractional GPU compute modes diff --git a/pkg/scheduler/actions/allocate/allocateFractionalGpu_test.go b/pkg/scheduler/actions/allocate/allocateFractionalGpu_test.go index dc87ae3ff..522c16837 100644 --- a/pkg/scheduler/actions/allocate/allocateFractionalGpu_test.go +++ b/pkg/scheduler/actions/allocate/allocateFractionalGpu_test.go @@ -9,13 +9,20 @@ import ( . "go.uber.org/mock/gomock" "gopkg.in/h2non/gock.v1" v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" "k8s.io/utils/ptr" kaiv1common "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1/common" + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" commonconstants "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/actions/allocate" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/actions/integration_tests/integration_tests_utils" + "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/pod_status" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/resource_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/conf" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/test_utils" @@ -42,6 +49,69 @@ func TestHandleFractionalGPUAllocation(t *testing.T) { } } +func TestFractionalGPUAllocationDoesNotUseGpuGroupWithDifferentComputeMode(t *testing.T) { + test_utils.InitTestingInfrastructure() + controller := NewController(t) + defer controller.Finish() + defer gock.Off() + + const gpuGroup = "time-slicing-group" + topology := test_utils.TestTopologyBasic{ + Name: "sm-sharing pod does not allocate on time-slicing gpu group", + Jobs: []*jobs_fake.TestJobBasic{ + { + Name: "running_job0", + RequiredGPUsPerTask: 0.5, + Priority: constants.PriorityTrainNumber, + QueueName: "queue0", + Tasks: []*tasks_fake.TestTaskBasic{ + { + NodeName: "node0", + GPUGroups: []string{gpuGroup}, + State: pod_status.Running, + }, + }, + }, + { + Name: "pending_job0", + RequiredGPUsPerTask: 0.5, + Priority: constants.PriorityTrainNumber, + QueueName: "queue0", + Tasks: []*tasks_fake.TestTaskBasic{ + { + State: pod_status.Pending, + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeSMSharing), + }, + }, + }, + }, + }, + Nodes: map[string]nodes_fake.TestNodeBasic{ + "node0": { + GPUs: 1, + }, + }, + Queues: []test_utils.TestQueueBasic{ + { + Name: "queue0", + DeservedGPUs: 1, + }, + }, + } + + ssn := test_utils.BuildSession(topology, controller) + addReservationPodToNodeForTest(ssn.ClusterInfo.Nodes["node0"], gpuGroup, schedulingv1alpha2.GPUComputeSharingModeTimeSlicing) + + allocate.New().Execute(ssn) + + pendingTask := ssn.ClusterInfo.PodGroupInfos["pending_job0"].GetAllPodsMap()[common_info.PodID("pending_job0-0")] + if pendingTask.Status != pod_status.Pending { + t.Fatalf("expected sm-sharing task to stay pending, got status %s on node %s with groups %v", + pendingTask.Status, pendingTask.NodeName, pendingTask.GPUGroupIDs()) + } +} + func TestFractionalGPUAllocationUsesNodeConditionOverride(t *testing.T) { test_utils.InitTestingInfrastructure() controller := NewController(t) @@ -111,6 +181,35 @@ func TestFractionalGPUAllocationUsesNodeConditionOverride(t *testing.T) { test_utils.MatchExpectedAndRealTasks(t, 0, topology, ssn) } +func addReservationPodToNodeForTest( + node *node_info.NodeInfo, gpuGroup string, mode schedulingv1alpha2.GPUComputeSharingMode, +) { + pod := &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + UID: types.UID("reservation-" + gpuGroup), + Name: commonconstants.GPUReservationPodPrefix + "-node0-test", + Namespace: commonconstants.DefaultResourceReservationName, + Labels: map[string]string{ + commonconstants.GPUGroup: gpuGroup, + }, + Annotations: map[string]string{ + commonconstants.PodGroupAnnotationForPod: "reservation", + commonconstants.GpuComputeSharingMode: string(mode), + }, + }, + Spec: v1.PodSpec{ + NodeName: node.Name, + Containers: []v1.Container{ + {Name: "reservation"}, + }, + }, + Status: v1.PodStatus{Phase: v1.PodRunning}, + } + task := pod_info.NewTaskInfo(pod, resource_info.NewResourceVectorMap()) + task.Status = pod_status.Running + node.PodInfos[task.UID] = task +} + func getFractionalGPUTestsMetadata() []integration_tests_utils.TestTopologyMetadata { return []integration_tests_utils.TestTopologyMetadata{ { diff --git a/pkg/scheduler/actions/integration_tests/integration_tests_utils/integration_tests_utils.go b/pkg/scheduler/actions/integration_tests/integration_tests_utils/integration_tests_utils.go index 88e56d05e..8f7afc6c8 100644 --- a/pkg/scheduler/actions/integration_tests/integration_tests_utils/integration_tests_utils.go +++ b/pkg/scheduler/actions/integration_tests/integration_tests_utils/integration_tests_utils.go @@ -110,7 +110,7 @@ func runSchedulerOneRound(testMetadata *TestTopologyMetadata, controller *Contro case pod_status.Releasing: if jobMetadata.DeleteJobInTest { taskMetadata.NodeName = task.NodeName - taskMetadata.GPUGroups = task.GPUGroups + taskMetadata.GPUGroups = task.GPUGroupIDs() taskMetadata.State = pod_status.Releasing } else { taskMetadata.NodeName = "" @@ -124,12 +124,12 @@ func runSchedulerOneRound(testMetadata *TestTopologyMetadata, controller *Contro case pod_status.Binding: taskMetadata.State = pod_status.Running taskMetadata.NodeName = task.NodeName - taskMetadata.GPUGroups = task.GPUGroups + taskMetadata.GPUGroups = task.GPUGroupIDs() default: taskMetadata.State = task.Status taskMetadata.NodeName = task.NodeName - taskMetadata.GPUGroups = task.GPUGroups + taskMetadata.GPUGroups = task.GPUGroupIDs() } } diff --git a/pkg/scheduler/actions/integration_tests/reclaim/reclaimFractional_test.go b/pkg/scheduler/actions/integration_tests/reclaim/reclaimFractional_test.go index 64eb5da98..f779bc573 100644 --- a/pkg/scheduler/actions/integration_tests/reclaim/reclaimFractional_test.go +++ b/pkg/scheduler/actions/integration_tests/reclaim/reclaimFractional_test.go @@ -6,8 +6,20 @@ package reclaim_test import ( "testing" + . "go.uber.org/mock/gomock" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" + commonconstants "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/actions/integration_tests/integration_tests_utils" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/actions/reclaim" + "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/pod_status" + "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/resource_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/test_utils" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/test_utils/jobs_fake" @@ -19,6 +31,128 @@ func TestReclaimFractionalIntegrationTest(t *testing.T) { integration_tests_utils.RunTests(t, getReclaimFractionalTestsMetadata()) } +func TestReclaimCanUseFullyEvictedFractionalGpuForDifferentComputeMode(t *testing.T) { + test_utils.InitTestingInfrastructure() + controller := NewController(t) + defer controller.Finish() + + const gpuGroup = "time-slicing-group" + topology := test_utils.TestTopologyBasic{ + Name: "reclaim all time-slicing fractions before using GPU for sm-sharing", + Jobs: []*jobs_fake.TestJobBasic{ + { + Name: "running_job0", + RequiredGPUsPerTask: 0.5, + Priority: constants.PriorityTrainNumber, + QueueName: "queue0", + Tasks: []*tasks_fake.TestTaskBasic{ + { + NodeName: "node0", + GPUGroups: []string{gpuGroup}, + State: pod_status.Running, + }, + { + NodeName: "node0", + GPUGroups: []string{gpuGroup}, + State: pod_status.Running, + }, + }, + }, + { + Name: "pending_job0", + RequiredGPUsPerTask: 0.5, + Priority: constants.PriorityTrainNumber, + QueueName: "queue1", + Tasks: []*tasks_fake.TestTaskBasic{ + { + State: pod_status.Pending, + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeSMSharing), + }, + }, + }, + }, + }, + Nodes: map[string]nodes_fake.TestNodeBasic{ + "node0": { + GPUs: 1, + }, + }, + Queues: []test_utils.TestQueueBasic{ + { + Name: "queue0", + DeservedGPUs: 0, + }, + { + Name: "queue1", + DeservedGPUs: 1, + }, + }, + Mocks: &test_utils.TestMock{ + CacheRequirements: &test_utils.CacheMocking{ + NumberOfCacheEvictions: 2, + NumberOfPipelineActions: 1, + }, + }, + } + + ssn := test_utils.BuildSession(topology, controller) + addReservationPodToNodeForReclaimTest(ssn.ClusterInfo.Nodes["node0"], gpuGroup, schedulingv1alpha2.GPUComputeSharingModeTimeSlicing) + + reclaim.New().Execute(ssn) + + runningTasks := ssn.ClusterInfo.PodGroupInfos["running_job0"].GetAllPodsMap() + for taskID, task := range runningTasks { + if task.Status != pod_status.Releasing { + t.Fatalf("expected time-slicing task %s to be releasing, got %s", taskID, task.Status) + } + } + + pendingTask := ssn.ClusterInfo.PodGroupInfos["pending_job0"].GetAllPodsMap()[common_info.PodID("pending_job0-0")] + if pendingTask.Status != pod_status.Pipelined { + t.Fatalf("expected sm-sharing task to be pipelined, got status %s", pendingTask.Status) + } + if pendingTask.NodeName != "node0" { + t.Fatalf("expected sm-sharing task on node0, got %s", pendingTask.NodeName) + } + gpuGroups := pendingTask.GPUGroupIDs() + if len(gpuGroups) != 1 || gpuGroups[0] == gpuGroup { + t.Fatalf("expected sm-sharing task to use a new gpu group, got %v", gpuGroups) + } + if pendingTask.FractionalGpuGroups[0].ComputeSharingMode != schedulingv1alpha2.GPUComputeSharingModeSMSharing { + t.Fatalf("expected sm-sharing mode, got %s", pendingTask.FractionalGpuGroups[0].ComputeSharingMode) + } +} + +func addReservationPodToNodeForReclaimTest( + node *node_info.NodeInfo, gpuGroup string, mode schedulingv1alpha2.GPUComputeSharingMode, +) { + pod := &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + UID: types.UID("reservation-" + gpuGroup), + Name: commonconstants.GPUReservationPodPrefix + "-node0-test", + Namespace: commonconstants.DefaultResourceReservationName, + Labels: map[string]string{ + commonconstants.GPUGroup: gpuGroup, + }, + Annotations: map[string]string{ + commonconstants.PodGroupAnnotationForPod: "reservation", + commonconstants.GpuComputeSharingMode: string(mode), + }, + }, + Spec: v1.PodSpec{ + NodeName: node.Name, + Containers: []v1.Container{ + {Name: "reservation"}, + }, + }, + Status: v1.PodStatus{Phase: v1.PodRunning}, + } + task := pod_info.NewTaskInfo(pod, resource_info.NewResourceVectorMap()) + task.Status = pod_status.Running + node.PodInfos[task.UID] = task +} + func getReclaimFractionalTestsMetadata() []integration_tests_utils.TestTopologyMetadata { return []integration_tests_utils.TestTopologyMetadata{ { diff --git a/pkg/scheduler/api/node_info/gpu_sharing_node_info.go b/pkg/scheduler/api/node_info/gpu_sharing_node_info.go index 4bf166b76..637b1dc15 100644 --- a/pkg/scheduler/api/node_info/gpu_sharing_node_info.go +++ b/pkg/scheduler/api/node_info/gpu_sharing_node_info.go @@ -6,7 +6,10 @@ package node_info import ( "fmt" "math" + "strings" + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" + commonconstants "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants" "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/resource_info" @@ -63,6 +66,53 @@ func (g *GpuSharingNodeInfo) Clone() *GpuSharingNodeInfo { return gpuSharingNodeInfo } +func (ni *NodeInfo) IsGpuGroupComputeSharingModeCompatible(gpuGroup string, task *pod_info.PodInfo) bool { + requestedMode := task.RequestedGPUComputeSharingMode() + if mode, found := ni.getReservationPodGpuGroupComputeSharingMode(gpuGroup); found { + return mode == requestedMode + } + if mode, found := ni.getGpuGroupComputeSharingMode(gpuGroup); found { + return mode == requestedMode + } + for _, fractionalGpuGroup := range task.FractionalGpuGroupsOrDefault() { + fractionalGpuGroup = fractionalGpuGroup.WithDefaults() + if fractionalGpuGroup.ID == gpuGroup { + return fractionalGpuGroup.ComputeSharingMode == requestedMode + } + } + return requestedMode == schedulingv1alpha2.GPUComputeSharingModeTimeSlicing +} + +func (ni *NodeInfo) getReservationPodGpuGroupComputeSharingMode( + gpuGroup string, +) (schedulingv1alpha2.GPUComputeSharingMode, bool) { + for _, task := range ni.PodInfos { + if task.Pod == nil || + !strings.HasPrefix(task.Pod.Name, commonconstants.GPUReservationPodPrefix) || + task.Pod.Labels[commonconstants.GPUGroup] != gpuGroup { + continue + } + if task.Pod.Annotations == nil { + return schedulingv1alpha2.GPUComputeSharingModeTimeSlicing, true + } + return schedulingv1alpha2.DefaultGPUComputeSharingMode( + schedulingv1alpha2.GPUComputeSharingMode(task.Pod.Annotations[commonconstants.GpuComputeSharingMode]), + ), true + } + return "", false +} + +func (ni *NodeInfo) getGpuGroupComputeSharingMode(gpuGroup string) (schedulingv1alpha2.GPUComputeSharingMode, bool) { + for _, task := range ni.PodInfos { + for _, fractionalGpuGroup := range task.FractionalGpuGroupsOrDefault() { + if fractionalGpuGroup.ID == gpuGroup { + return fractionalGpuGroup.WithDefaults().ComputeSharingMode, true + } + } + } + return "", false +} + /************* All - Shared Tasks *************/ func getAcceptedTaskResourceVectorWithoutSharedGPU(task *pod_info.PodInfo, vectorMap *resource_info.ResourceVectorMap) resource_info.ResourceVector { @@ -81,7 +131,7 @@ func (ni *NodeInfo) addSharedGPUTaskResources(task *pod_info.PodInfo) { log.InfraLogger.V(7).Infof("About to add shared podsInfo: <%v/%v>, status: <%v>, node: <%+v>", task.Namespace, task.Name, task.Status, ni) - for _, gpuGroup := range task.GPUGroups { + for _, gpuGroup := range task.GPUGroupIDs() { ni.addSharedGPUTaskResourcesPerPodGroup(task, gpuGroup) } @@ -93,7 +143,7 @@ func (ni *NodeInfo) addSharedGPUTaskResourcesPerPodGroup(task *pod_info.PodInfo, log.InfraLogger.V(7).Infof( "About to add shared podsInfo: <%v/%v>, gpuGroup: <%v> "+ "releasingSharedGPU: <%v> AllocatedSharedGPUsMemory <%v>, UsedSharedGPUsMemory: <%v>", - task.Namespace, task.Name, task.GPUGroups, + task.Namespace, task.Name, task.GPUGroupIDs(), ni.ReleasingSharedGPUsMemory[gpuGroup], ni.AllocatedSharedGPUsMemory[gpuGroup], ni.UsedSharedGPUsMemory[gpuGroup]) @@ -141,7 +191,7 @@ func (ni *NodeInfo) addSharedGPUTaskResourcesPerPodGroup(task *pod_info.PodInfo, log.InfraLogger.V(8).Infof( "Added shared podsInfo: <%v/%v>, gpuGroup: <%v> "+ "releasingSharedGPU: <%v> AllocatedSharedGPUsMemory <%v>, UsedSharedGPUsMemory: <%v>", - task.Namespace, task.Name, task.GPUGroups, + task.Namespace, task.Name, task.GPUGroupIDs(), ni.ReleasingSharedGPUsMemory[gpuGroup], ni.AllocatedSharedGPUsMemory[gpuGroup], ni.UsedSharedGPUsMemory[gpuGroup]) } @@ -153,9 +203,9 @@ func (ni *NodeInfo) removeSharedTaskResources(task *pod_info.PodInfo) { log.InfraLogger.V(7).Infof( "About to remove shared podsInfo: <%v/%v>, status: <%v>, node: <%+v>", - task.Namespace, task.Name, task.GPUGroups, ni) + task.Namespace, task.Name, task.GPUGroupIDs(), ni) - for _, gpuGroup := range task.GPUGroups { + for _, gpuGroup := range task.GPUGroupIDs() { ni.removeSharedTaskResourcesPerPodGroup(task, gpuGroup) } @@ -168,7 +218,7 @@ func (ni *NodeInfo) removeSharedTaskResourcesPerPodGroup(task *pod_info.PodInfo, log.InfraLogger.V(7).Infof( "About to remove shared podsInfo: <%v/%v>, gpuGroup: <%v> "+ "releasingSharedGPU: <%v> AllocatedSharedGPUsMemory <%v>, UsedSharedGPUsMemory: <%v>", - task.Namespace, task.Name, task.GPUGroups, + task.Namespace, task.Name, task.GPUGroupIDs(), ni.ReleasingSharedGPUsMemory[gpuGroup], ni.AllocatedSharedGPUsMemory[gpuGroup], ni.UsedSharedGPUsMemory[gpuGroup]) @@ -231,7 +281,7 @@ func (ni *NodeInfo) removeSharedTaskResourcesPerPodGroup(task *pod_info.PodInfo, log.InfraLogger.V(8).Infof( "Removed shared podsInfo: <%v/%v>, gpuGroup: <%v> "+ "releasingSharedGPU: <%v> AllocatedSharedGPUsMemory <%v>, UsedSharedGPUsMemory: <%v>", - task.Namespace, task.Name, task.GPUGroups, + task.Namespace, task.Name, task.GPUGroupIDs(), ni.ReleasingSharedGPUsMemory[gpuGroup], ni.AllocatedSharedGPUsMemory[gpuGroup], ni.UsedSharedGPUsMemory[gpuGroup]) } diff --git a/pkg/scheduler/api/node_info/node_info_test.go b/pkg/scheduler/api/node_info/node_info_test.go index e70f0389c..4f1080d4e 100644 --- a/pkg/scheduler/api/node_info/node_info_test.go +++ b/pkg/scheduler/api/node_info/node_info_test.go @@ -35,6 +35,7 @@ import ( "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" commonconstants "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants" "github.com/kai-scheduler/KAI-scheduler/pkg/common/resources" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/common_info" @@ -657,7 +658,7 @@ func TestAddRemovePods(t *testing.T) { for _, podInfoMetaData := range test.podsInfoMetadata { pi := pod_info.NewTaskInfo(podInfoMetaData.pod, vectorMap) pi.Status = podInfoMetaData.status - pi.GPUGroups = podInfoMetaData.gpuGroups + pi.SetGPUGroupIDs(podInfoMetaData.gpuGroups) podsInfo = append(podsInfo, pi) } @@ -1359,6 +1360,89 @@ func TestNodeInfo_GetSumOfReleasingGPUs(t *testing.T) { } } +func TestIsGpuGroupComputeSharingModeCompatible_CurrentTaskNewGroup(t *testing.T) { + gpuGroup := "new-sm-sharing-group" + node := &NodeInfo{ + PodInfos: map[common_info.PodID]*pod_info.PodInfo{}, + GpuSharingNodeInfo: GpuSharingNodeInfo{ + UsedSharedGPUsMemory: map[string]int64{ + gpuGroup: 0, + }, + }, + } + task := &pod_info.PodInfo{ + Pod: &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeSMSharing), + }, + }, + }, + FractionalGpuGroups: []schedulingv1alpha2.FractionalGpuGroup{ + { + ID: gpuGroup, + ComputeSharingMode: schedulingv1alpha2.GPUComputeSharingModeSMSharing, + }, + }, + } + task.SetGPUGroupIDs([]string{gpuGroup}) + + assert.True(t, node.IsGpuGroupComputeSharingModeCompatible(gpuGroup, task)) +} + +func TestIsGpuGroupComputeSharingModeCompatible_ReservationPodIsSourceOfTruth(t *testing.T) { + gpuGroup := "reserved-group" + node := &NodeInfo{ + PodInfos: map[common_info.PodID]*pod_info.PodInfo{ + "reservation": { + Pod: &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: commonconstants.GPUReservationPodPrefix + "-node-a-abcde", + Labels: map[string]string{ + commonconstants.GPUGroup: gpuGroup, + }, + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeTimeSlicing), + }, + }, + }, + }, + "workload": { + Pod: &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeSMSharing), + }, + }, + }, + FractionalGpuGroups: []schedulingv1alpha2.FractionalGpuGroup{ + { + ID: gpuGroup, + ComputeSharingMode: schedulingv1alpha2.GPUComputeSharingModeSMSharing, + }, + }, + }, + }, + } + task := &pod_info.PodInfo{ + Pod: &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Annotations: map[string]string{ + commonconstants.GpuComputeSharingMode: string(schedulingv1alpha2.GPUComputeSharingModeSMSharing), + }, + }, + }, + FractionalGpuGroups: []schedulingv1alpha2.FractionalGpuGroup{ + { + ID: gpuGroup, + ComputeSharingMode: schedulingv1alpha2.GPUComputeSharingModeSMSharing, + }, + }, + } + + assert.False(t, node.IsGpuGroupComputeSharingModeCompatible(gpuGroup, task)) +} + func createPod(namespace, name string, options podCreationOptions) *pod_info.PodInfo { pod := &v1.Pod{ ObjectMeta: metav1.ObjectMeta{ @@ -1397,7 +1481,7 @@ func createPod(namespace, name string, options podCreationOptions) *pod_info.Pod } task := pod_info.NewTaskInfo(pod, resource_info.NewResourceVectorMap()) - task.GPUGroups = []string{options.gpuGroup} + task.SetGPUGroupIDs([]string{options.gpuGroup}) return task } diff --git a/pkg/scheduler/api/pod_info/pod_info.go b/pkg/scheduler/api/pod_info/pod_info.go index 7e3b0866c..4840915a0 100644 --- a/pkg/scheduler/api/pod_info/pod_info.go +++ b/pkg/scheduler/api/pod_info/pod_info.go @@ -95,7 +95,7 @@ type PodInfo struct { schedulingConstraintsSignature common_info.SchedulingConstraintsSignature - GPUGroups []string + FractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup NUMAPlacement NUMAPlacement @@ -215,7 +215,6 @@ func NewTaskInfo(pod *v1.Pod, vectorMap *resource_info.ResourceVectorMap, opts . ResReqVector: initResreq.ToVector(vectorMap), AcceptedResourceVector: resource_info.NewResourceVector(vectorMap), VectorMap: vectorMap, - GPUGroups: []string{}, ResourceRequestType: RequestTypeRegular, ResourceReceivedType: ReceivedTypeNone, BindRequest: options.BindRequest, @@ -302,8 +301,8 @@ func (pi *PodInfo) Clone() *PodInfo { ResReqVector: resReqVectorClone, AcceptedResourceVector: acceptedResourceVectorClone, VectorMap: pi.VectorMap, - GPUGroups: pi.GPUGroups, NUMAPlacement: pi.NUMAPlacement.Clone(), + FractionalGpuGroups: pi.FractionalGpuGroups, ResourceClaimInfo: pi.ResourceClaimInfo.Clone(), ExtendedResourceClaim: pi.ExtendedResourceClaim, ResourceRequestType: pi.ResourceRequestType, @@ -315,6 +314,53 @@ func (pi *PodInfo) Clone() *PodInfo { } } +func (pi *PodInfo) SetFractionalGpuGroups(fractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup) { + if len(fractionalGpuGroups) == 0 { + pi.FractionalGpuGroups = nil + return + } + pi.FractionalGpuGroups = fractionalGpuGroups +} + +func (pi *PodInfo) SetGPUGroupIDs(gpuGroups []string) { + pi.SetFractionalGpuGroups(schedulingv1alpha2.NewFractionalGpuGroups( + gpuGroups, + pi.RequestedGPUComputeSharingMode(), + )) +} + +func (pi *PodInfo) GPUGroupIDs() []string { + fractionalGpuGroups := pi.FractionalGpuGroupsOrDefault() + if len(fractionalGpuGroups) == 0 { + return nil + } + gpuGroups := make([]string, 0, len(fractionalGpuGroups)) + for _, fractionalGpuGroup := range fractionalGpuGroups { + gpuGroups = append(gpuGroups, fractionalGpuGroup.ID) + } + return gpuGroups +} + +func (pi *PodInfo) RequestedGPUComputeSharingMode() schedulingv1alpha2.GPUComputeSharingMode { + if pi.Pod == nil || pi.Pod.Annotations == nil { + return schedulingv1alpha2.GPUComputeSharingModeTimeSlicing + } + return schedulingv1alpha2.DefaultGPUComputeSharingMode( + schedulingv1alpha2.GPUComputeSharingMode(pi.Pod.Annotations[commonconstants.GpuComputeSharingMode]), + ) +} + +func (pi *PodInfo) FractionalGpuGroupsOrDefault() []schedulingv1alpha2.FractionalGpuGroup { + if len(pi.FractionalGpuGroups) > 0 { + fractionalGpuGroups := make([]schedulingv1alpha2.FractionalGpuGroup, 0, len(pi.FractionalGpuGroups)) + for _, fractionalGpuGroup := range pi.FractionalGpuGroups { + fractionalGpuGroups = append(fractionalGpuGroups, fractionalGpuGroup.WithDefaults()) + } + return fractionalGpuGroups + } + return nil +} + func (pi PodInfo) String() string { return fmt.Sprintf("Pod (%v:%v/%v): job %v, status %v, resreq %v, gpu %v", pi.UID, pi.Namespace, pi.Name, pi.Job, pi.Status, pi.ResReqVector, pi.GpuRequirement.GpusAsString()) @@ -503,10 +549,13 @@ func getTaskStatus(pod *v1.Pod, bindRequest *bindrequest_info.BindRequestInfo, s } func (pi *PodInfo) updatePodAdditionalFields(bindRequest *bindrequest_info.BindRequestInfo, draPodClaims ...*resourceapi.ResourceClaim) { - if bindRequest != nil && len(bindRequest.BindRequest.Spec.SelectedGPUGroups) > 0 { - pi.GPUGroups = bindRequest.BindRequest.Spec.SelectedGPUGroups + if bindRequest != nil && len(bindRequest.BindRequest.Spec.SelectedFractionalGpuGroupsOrDefault()) > 0 { + pi.SetFractionalGpuGroups(bindRequest.BindRequest.Spec.SelectedFractionalGpuGroupsOrDefault()) } else { - pi.GPUGroups = resources.GetGpuGroups(pi.Pod) + pi.SetFractionalGpuGroups(schedulingv1alpha2.NewFractionalGpuGroups( + resources.GetGpuGroups(pi.Pod), + pi.RequestedGPUComputeSharingMode(), + )) } if bindRequest != nil && len(bindRequest.BindRequest.Spec.ReceivedResourceType) > 0 { diff --git a/pkg/scheduler/api/pod_info/pod_info_test.go b/pkg/scheduler/api/pod_info/pod_info_test.go index 9f4c0642e..610b6b054 100644 --- a/pkg/scheduler/api/pod_info/pod_info_test.go +++ b/pkg/scheduler/api/pod_info/pod_info_test.go @@ -596,7 +596,6 @@ func TestPodInfo_updatePodAdditionalFields(t *testing.T) { Status: tt.fields.Status, Pod: tt.fields.Pod, GpuRequirement: *resource_info.NewGpuResourceRequirement(), - GPUGroups: make([]string, 0), VectorMap: vectorMap, } pi.updatePodAdditionalFields(tt.fields.bindingRequest) @@ -610,9 +609,9 @@ func TestPodInfo_updatePodAdditionalFields(t *testing.T) { tt.name, tt.expected.AcceptedGpuRequirement, pi.AcceptedGpuRequirement) } assert.Equal(t, string(pi.ResourceRequestType), tt.expected.ResourceRequestType) - if !reflect.DeepEqual(pi.GPUGroups, tt.expected.GPUGroups) { + if !reflect.DeepEqual(pi.GPUGroupIDs(), tt.expected.GPUGroups) { t.Errorf("case (%s) failed: GPUGroups \n expected %v, \n got: %v \n", - tt.name, tt.expected.GPUGroups, pi.GPUGroups) + tt.name, tt.expected.GPUGroups, pi.GPUGroupIDs()) } }) } diff --git a/pkg/scheduler/cache/cache.go b/pkg/scheduler/cache/cache.go index 796eda86e..1bdc4897d 100644 --- a/pkg/scheduler/cache/cache.go +++ b/pkg/scheduler/cache/cache.go @@ -332,7 +332,7 @@ func (sc *SchedulerCache) Bind(taskInfo *pod_info.PodInfo, hostname string, bind log.InfraLogger.V(3).Infof( "Creating bind request for task <%v/%v> to node <%v> gpuGroup: <%v>, requires: <%v> GPUs", - taskInfo.Namespace, taskInfo.Name, hostname, taskInfo.GPUGroups, taskInfo.ResReqVector) + taskInfo.Namespace, taskInfo.Name, hostname, taskInfo.GPUGroupIDs(), taskInfo.ResReqVector) if bindRequestError := sc.createBindRequest(taskInfo, hostname, bindRequestAnnotations, predictedNUMAZones); bindRequestError != nil { return sc.StatusUpdater.Bound(taskInfo.Pod, hostname, bindRequestError, sc.getNodPoolName()) } @@ -372,10 +372,11 @@ func (sc *SchedulerCache) createBindRequest(podInfo *pod_info.PodInfo, nodeName Labels: labels, }, Spec: schedulingv1alpha2.BindRequestSpec{ - PodName: podInfo.Name, - SelectedNode: nodeName, - SelectedGPUGroups: podInfo.GPUGroups, - ReceivedResourceType: string(podInfo.ResourceReceivedType), + PodName: podInfo.Name, + SelectedNode: nodeName, + SelectedGPUGroups: podInfo.GPUGroupIDs(), + SelectedFractionalGpuGroups: podInfo.FractionalGpuGroupsOrDefault(), + ReceivedResourceType: string(podInfo.ResourceReceivedType), ReceivedGPU: &schedulingv1alpha2.ReceivedGPU{ Count: int(podInfo.AcceptedGpuRequirement.GetNumOfGpuDevices()), Portion: fmt.Sprintf("%.2f", podInfo.AcceptedGpuRequirement.GpuFractionalPortion()), diff --git a/pkg/scheduler/framework/operations.go b/pkg/scheduler/framework/operations.go index 5da0b9960..c206fc72c 100644 --- a/pkg/scheduler/framework/operations.go +++ b/pkg/scheduler/framework/operations.go @@ -20,6 +20,7 @@ limitations under the License. package framework import ( + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/bindrequest_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/eviction_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/node_info" @@ -47,14 +48,14 @@ type Operation interface { type ReverseOperation func() error type evictOperation struct { - taskInfo *pod_info.PodInfo - previousStatus pod_status.PodStatus - previousNode *node_info.NodeInfo - previousGpuGroups []string - previousNumaPlacement pod_info.NUMAPlacement - message string - evictionMetadata eviction_info.EvictionMetadata - reverseOperation ReverseOperation + taskInfo *pod_info.PodInfo + previousStatus pod_status.PodStatus + previousNode *node_info.NodeInfo + previousFractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup + previousNumaPlacement pod_info.NUMAPlacement + message string + evictionMetadata eviction_info.EvictionMetadata + reverseOperation ReverseOperation } func (op evictOperation) Name() string { @@ -88,15 +89,15 @@ func (op allocateOperation) Reverse() error { } type pipelineOperation struct { - taskInfo *pod_info.PodInfo - previousStatus pod_status.PodStatus - previousNode string - previousGpuGroups []string - previousNumaPlacement pod_info.NUMAPlacement - previousResourceClaimInfo bindrequest_info.ResourceClaimInfo - nextNode string - message string - reverseOperation ReverseOperation + taskInfo *pod_info.PodInfo + previousStatus pod_status.PodStatus + previousNode string + previousFractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup + previousNumaPlacement pod_info.NUMAPlacement + previousResourceClaimInfo bindrequest_info.ResourceClaimInfo + nextNode string + message string + reverseOperation ReverseOperation } func (op pipelineOperation) Name() string { diff --git a/pkg/scheduler/framework/session.go b/pkg/scheduler/framework/session.go index 099a05548..1fbeec2cf 100644 --- a/pkg/scheduler/framework/session.go +++ b/pkg/scheduler/framework/session.go @@ -193,7 +193,8 @@ func (ssn *Session) FittingGPUs(node *node_info.NodeInfo, pod *pod_info.PodInfo) func filterGpusByEnoughResources(node *node_info.NodeInfo, pod *pod_info.PodInfo) []string { filteredGPUs := []string{} for gpuIdx := range node.UsedSharedGPUsMemory { - if node.IsTaskFitOnGpuGroup(&pod.GpuRequirement, gpuIdx) { + if node.IsTaskFitOnGpuGroup(&pod.GpuRequirement, gpuIdx) && + node.IsGpuGroupComputeSharingModeCompatible(gpuIdx, pod) { filteredGPUs = append(filteredGPUs, gpuIdx) } } diff --git a/pkg/scheduler/framework/statement.go b/pkg/scheduler/framework/statement.go index 79bf08be5..bdf9f141b 100644 --- a/pkg/scheduler/framework/statement.go +++ b/pkg/scheduler/framework/statement.go @@ -25,6 +25,7 @@ import ( "golang.org/x/exp/slices" + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/bindrequest_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/common_info" "github.com/kai-scheduler/KAI-scheduler/pkg/scheduler/api/eviction_info" @@ -78,7 +79,7 @@ func (s *Statement) Evict(reclaimeeTask *pod_info.PodInfo, message string, } previousStatus := reclaimeeTask.Status - previousGpuGroup := reclaimeeTask.GPUGroups + previousFractionalGpuGroups := cloneFractionalGpuGroups(reclaimeeTask.FractionalGpuGroups) previousNumaPlacement := reclaimeeTask.NUMAPlacement.Clone() previousIsVirtualStatus := reclaimeeTask.IsVirtualStatus var previousResourceClaimInfo bindrequest_info.ResourceClaimInfo @@ -107,15 +108,17 @@ func (s *Statement) Evict(reclaimeeTask *pod_info.PodInfo, message string, s.operations = append(s.operations, evictOperation{ - taskInfo: reclaimeeTask, - previousStatus: previousStatus, - previousNode: node, - previousGpuGroups: previousGpuGroup, - previousNumaPlacement: previousNumaPlacement, - message: message, - evictionMetadata: evictionMetadata, + taskInfo: reclaimeeTask, + previousStatus: previousStatus, + previousNode: node, + previousFractionalGpuGroups: previousFractionalGpuGroups, + previousNumaPlacement: previousNumaPlacement, + message: message, + evictionMetadata: evictionMetadata, reverseOperation: func() error { - return s.unevict(reclaimeeTask, previousStatus, node, previousGpuGroup, previousNumaPlacement, previousResourceClaimInfo, previousIsVirtualStatus) + return s.unevict( + reclaimeeTask, previousStatus, node, previousFractionalGpuGroups, + previousNumaPlacement, previousResourceClaimInfo, previousIsVirtualStatus) }, }, ) @@ -135,14 +138,15 @@ func (s *Statement) commitEvict(reclaimee *pod_info.PodInfo, evictOp evictOperat } previousStatus := reclaimee.Status - previousGpuGroup := reclaimee.GPUGroups + previousFractionalGpuGroups := cloneFractionalGpuGroups(reclaimee.FractionalGpuGroups) previousNumaPlacement := reclaimee.NUMAPlacement.Clone() previousResourceClaimInfo := reclaimee.ResourceClaimInfo previousIsVirtualStatus := reclaimee.IsVirtualStatus if err := s.ssn.Cache.Evict(reclaimee.Pod, reclaimeePodGroup, evictOp.evictionMetadata, evictOp.message); err != nil { log.InfraLogger.Errorf("Failed to evict task <%v/%v>: %v.", reclaimee.Namespace, reclaimee.Name, err) - if e := s.unevict(reclaimee, previousStatus, evictOp.previousNode, previousGpuGroup, previousNumaPlacement, previousResourceClaimInfo, - previousIsVirtualStatus); e != nil { + if e := s.unevict( + reclaimee, previousStatus, evictOp.previousNode, previousFractionalGpuGroups, + previousNumaPlacement, previousResourceClaimInfo, previousIsVirtualStatus); e != nil { log.InfraLogger.Errorf("Failed to un-evict task <%v/%v>: %v.", reclaimee.Namespace, reclaimee.Name, e) } @@ -158,7 +162,8 @@ func (s *Statement) commitEvict(reclaimee *pod_info.PodInfo, evictOp evictOperat func (s *Statement) unevict( reclaimee *pod_info.PodInfo, previousStatus pod_status.PodStatus, node *node_info.NodeInfo, - previousGpuGroups []string, previousNumaPlacement pod_info.NUMAPlacement, + previousFractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup, + previousNumaPlacement pod_info.NUMAPlacement, previousResourceClaimInfo bindrequest_info.ResourceClaimInfo, previousIsVirtualStatus bool) error { // Update status in session job, found := s.ssn.ClusterInfo.PodGroupInfos[reclaimee.Job] @@ -171,7 +176,7 @@ func (s *Statement) unevict( log.InfraLogger.Errorf("Failed to find Job <%s> in Session <%s> index when binding.", reclaimee.Job, s.sessionID) } - reclaimee.GPUGroups = previousGpuGroups + reclaimee.FractionalGpuGroups = previousFractionalGpuGroups reclaimee.NUMAPlacement = previousNumaPlacement.Clone() reclaimee.IsVirtualStatus = previousIsVirtualStatus reclaimee.ResourceClaimInfo = previousResourceClaimInfo.Clone() @@ -217,8 +222,10 @@ func (s *Statement) Pipeline(task *pod_info.PodInfo, hostname string, updateTask gpuPlacementChanged := false numaPlacementChanged := false if foundOnNode { - gpuPlacementChanged = len(task.GPUGroups) > 0 && task.IsSharedGPUAllocation() && - !slices.Equal(task.GPUGroups, []string{"-1"}) && !slices.Equal(task.GPUGroups, taskOnNode.GPUGroups) + taskGPUGroups := task.GPUGroupIDs() + taskOnNodeGPUGroups := taskOnNode.GPUGroupIDs() + gpuPlacementChanged = len(taskGPUGroups) > 0 && task.IsSharedGPUAllocation() && + !slices.Equal(taskGPUGroups, []string{"-1"}) && !slices.Equal(taskGPUGroups, taskOnNodeGPUGroups) numaPlacementChanged = len(task.NUMAPlacement) > 0 && !task.NUMAPlacement.Equal(taskOnNode.NUMAPlacement) } @@ -227,7 +234,7 @@ func (s *Statement) Pipeline(task *pod_info.PodInfo, hostname string, updateTask // If task already exist on the node, and we didn't ask to update if on the node, // and there is no special reason we should update on the node, then we need to unevict instead of pipelining it. if foundOnNode && !updateTaskIfExistsOnNode && !gpuPlacementChanged && !numaPlacementChanged { - task.GPUGroups = taskOnNode.GPUGroups + task.FractionalGpuGroups = cloneFractionalGpuGroups(taskOnNode.FractionalGpuGroups) task.NUMAPlacement = taskOnNode.NUMAPlacement.Clone() log.InfraLogger.V(6).Infof("Task: <%v/%v> already exists on node: <%v>, unevicting it", task.Namespace, task.Name, hostname) if err := s.Unevict(task); err != nil { @@ -246,8 +253,7 @@ func (s *Statement) Pipeline(task *pod_info.PodInfo, hostname string, updateTask previousNode := task.NodeName task.NodeName = hostname - previousGpuGroup := task.GPUGroups - + previousFractionalGpuGroups := cloneFractionalGpuGroups(task.FractionalGpuGroups) previousNumaPlacement := task.NUMAPlacement.Clone() if numaPlacementChanged { previousNumaPlacement = taskOnNode.NUMAPlacement.Clone() @@ -261,8 +267,8 @@ func (s *Statement) Pipeline(task *pod_info.PodInfo, hostname string, updateTask if gpuPlacementChanged { log.InfraLogger.V(6).Infof( "Task: <%v/%v> already exists on node: <%v> on gpu index of: <%v>, moving it to index: <%v>", - task.Namespace, task.Name, hostname, taskOnNode.GPUGroups, task.GPUGroups) - previousGpuGroup = taskOnNode.GPUGroups + task.Namespace, task.Name, hostname, taskOnNode.GPUGroupIDs(), task.GPUGroupIDs()) + previousFractionalGpuGroups = cloneFractionalGpuGroups(taskOnNode.FractionalGpuGroups) if err := node.ConsolidateSharedPodInfoToDifferentGPU(task); err != nil { log.InfraLogger.Errorf("Failed to unevict task <%v/%v> to node <%v> in Session <%v>: %v", task.Namespace, task.Name, hostname, s.sessionID, err) @@ -291,23 +297,25 @@ func (s *Statement) Pipeline(task *pod_info.PodInfo, hostname string, updateTask } s.operations = append(s.operations, pipelineOperation{ - taskInfo: task, - previousStatus: previousStatus, - previousNode: previousNode, - previousGpuGroups: previousGpuGroup, - previousNumaPlacement: previousNumaPlacement, - previousResourceClaimInfo: previousResourceClaimInfo, - nextNode: hostname, - message: fmt.Sprintf("Pod %s/%s was pipelined to node %s", task.Namespace, task.Name, node.Name), + taskInfo: task, + previousStatus: previousStatus, + previousNode: previousNode, + previousFractionalGpuGroups: previousFractionalGpuGroups, + previousNumaPlacement: previousNumaPlacement, + previousResourceClaimInfo: previousResourceClaimInfo, + nextNode: hostname, + message: fmt.Sprintf("Pod %s/%s was pipelined to node %s", task.Namespace, task.Name, node.Name), reverseOperation: func() error { - return s.unpipeline(task, previousNode, previousStatus, previousGpuGroup, previousNumaPlacement, previousResourceClaimInfo, previousIsVirtualStatus) + return s.unpipeline( + task, previousNode, previousStatus, previousFractionalGpuGroups, + previousNumaPlacement, previousResourceClaimInfo, previousIsVirtualStatus) }, }) task.IsVirtualStatus = true log.InfraLogger.V(6).Infof( "Statement pipelined task: <%v/%v> to node: <%v>, gpuGroup: <%v>", - task.Namespace, task.Name, hostname, task.GPUGroups) + task.Namespace, task.Name, hostname, task.GPUGroupIDs()) return nil } @@ -395,7 +403,7 @@ func (s *Statement) commitAllocate(task *pod_info.PodInfo) error { }() if task.IsFractionAllocation() { - for _, gpuGroup := range task.GPUGroups { + for _, gpuGroup := range task.GPUGroupIDs() { if _, found := node.UsedSharedGPUsMemory[gpuGroup]; !found { node.UsedSharedGPUsMemory[gpuGroup] = 0 } @@ -456,7 +464,8 @@ func (s *Statement) commitPipeline(task *pod_info.PodInfo, message string) { } func (s *Statement) unpipeline( - task *pod_info.PodInfo, previousNode string, previousStatus pod_status.PodStatus, previousGpuGroups []string, + task *pod_info.PodInfo, previousNode string, previousStatus pod_status.PodStatus, + previousFractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup, previousNumaPlacement pod_info.NUMAPlacement, previousResourceClaimInfo bindrequest_info.ResourceClaimInfo, previousIsVirtualStatus bool) error { @@ -474,7 +483,7 @@ func (s *Statement) unpipeline( } hostname := task.NodeName - task.GPUGroups = previousGpuGroups + task.FractionalGpuGroups = previousFractionalGpuGroups task.ResourceClaimInfo = previousResourceClaimInfo.Clone() task.IsVirtualStatus = previousIsVirtualStatus @@ -692,3 +701,14 @@ func (s *Statement) operationValid(i int) bool { } return true } + +func cloneFractionalGpuGroups( + fractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup, +) []schedulingv1alpha2.FractionalGpuGroup { + if len(fractionalGpuGroups) == 0 { + return nil + } + clone := make([]schedulingv1alpha2.FractionalGpuGroup, len(fractionalGpuGroups)) + copy(clone, fractionalGpuGroups) + return clone +} diff --git a/pkg/scheduler/framework/statement_test.go b/pkg/scheduler/framework/statement_test.go index abe6b89df..db0e83df2 100644 --- a/pkg/scheduler/framework/statement_test.go +++ b/pkg/scheduler/framework/statement_test.go @@ -142,7 +142,7 @@ func TestStatement_Evict_Unevict(t *testing.T) { assert.Equal(t, actualTask.Status, originalTask.Status) assert.Equal(t, actualTask.GpuRequirement, originalTask.GpuRequirement) assert.Equal(t, actualTask.ResReqVector, originalTask.ResReqVector) - assert.Equal(t, actualTask.GPUGroups, originalTask.GPUGroups) + assert.Equal(t, actualTask.GPUGroupIDs(), originalTask.GPUGroupIDs()) actualJob := ssn.ClusterInfo.PodGroupInfos[tt.args.jobName] assert.Equal(t, originalJob.AllocatedVector, actualJob.AllocatedVector) @@ -637,7 +637,7 @@ func TestStatement_Pipeline_Unpipeline(t *testing.T) { assert.Equal(t, actualTask.Status, originalPipelineTask.Status) assert.Equal(t, actualTask.GpuRequirement, originalPipelineTask.GpuRequirement) assert.Equal(t, actualTask.ResReqVector, originalPipelineTask.ResReqVector) - assert.Equal(t, actualTask.GPUGroups, originalPipelineTask.GPUGroups) + assert.Equal(t, actualTask.GPUGroupIDs(), originalPipelineTask.GPUGroupIDs()) actualPipelinedJob := ssn.ClusterInfo.PodGroupInfos[tt.args.jobName] assert.Equal(t, originalPipelineJob.AllocatedVector, actualPipelinedJob.AllocatedVector) @@ -986,7 +986,7 @@ func TestStatement_Allocate_Unallocate(t *testing.T) { assert.Equal(t, actualAllocatedTask.Status, originalAllocateTask.Status) assert.Equal(t, actualAllocatedTask.GpuRequirement, originalAllocateTask.GpuRequirement) assert.Equal(t, actualAllocatedTask.ResReqVector, originalAllocateTask.ResReqVector) - assert.Equal(t, actualAllocatedTask.GPUGroups, originalAllocateTask.GPUGroups) + assert.Equal(t, actualAllocatedTask.GPUGroupIDs(), originalAllocateTask.GPUGroupIDs()) actualAllocatedJob := ssn.ClusterInfo.PodGroupInfos[tt.args.jobName] assert.Equal(t, originalAllocateJob.AllocatedVector, actualAllocatedJob.AllocatedVector) diff --git a/pkg/scheduler/gpu_sharing/gpuSharing.go b/pkg/scheduler/gpu_sharing/gpuSharing.go index c56445452..0365921f9 100644 --- a/pkg/scheduler/gpu_sharing/gpuSharing.go +++ b/pkg/scheduler/gpu_sharing/gpuSharing.go @@ -6,6 +6,7 @@ package gpu_sharing import ( "k8s.io/apimachinery/pkg/util/uuid" + schedulingv1alpha2 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/scheduling/v1alpha2" "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/framework" @@ -13,8 +14,19 @@ import ( ) type nodeGpuForSharing struct { - Groups []string - IsReleasing bool + FractionalGpuGroups []schedulingv1alpha2.FractionalGpuGroup + IsReleasing bool +} + +func (n *nodeGpuForSharing) GroupIDs() []string { + if n == nil { + return nil + } + groups := make([]string, 0, len(n.FractionalGpuGroups)) + for _, fractionalGpuGroup := range n.FractionalGpuGroups { + groups = append(groups, fractionalGpuGroup.ID) + } + return groups } func AllocateFractionalGPUTaskToNode(ssn *framework.Session, stmt *framework.Statement, pod *pod_info.PodInfo, @@ -25,12 +37,12 @@ func AllocateFractionalGPUTaskToNode(ssn *framework.Session, stmt *framework.Sta return false } - pod.GPUGroups = gpuForSharing.Groups + pod.SetFractionalGpuGroups(gpuForSharing.FractionalGpuGroups) isPipelineOnly = isPipelineOnly || gpuForSharing.IsReleasing success := allocateSharedGPUTask(ssn, stmt, node, pod, isPipelineOnly) if !success { - pod.GPUGroups = nil + pod.FractionalGpuGroups = nil } return success } @@ -39,8 +51,8 @@ func GetNodePreferableGpuForSharing(fittingGPUsOnNode []string, node *node_info. isPipelineOnly bool) *nodeGpuForSharing { nodeGpusSharing := &nodeGpuForSharing{ - Groups: []string{}, - IsReleasing: false, + FractionalGpuGroups: []schedulingv1alpha2.FractionalGpuGroup{}, + IsReleasing: false, } deviceCounts := pod.GpuRequirement.GetNumOfGpuDevices() @@ -49,17 +61,24 @@ func GetNodePreferableGpuForSharing(fittingGPUsOnNode []string, node *node_info. if wholeGpuForSharing := findGpuForSharingOnNode(pod, node, isPipelineOnly); wholeGpuForSharing != nil { nodeGpusSharing.IsReleasing = nodeGpusSharing.IsReleasing || wholeGpuForSharing.IsReleasing - nodeGpusSharing.Groups = append(nodeGpusSharing.Groups, wholeGpuForSharing.Groups...) + nodeGpusSharing.FractionalGpuGroups = append( + nodeGpusSharing.FractionalGpuGroups, wholeGpuForSharing.FractionalGpuGroups...) } } else { nodeGpusSharing.IsReleasing = nodeGpusSharing.IsReleasing || !node.EnoughIdleResourcesOnGpu(&pod.GpuRequirement, gpuIdx) || !node.IsTaskAllocatable(pod) - nodeGpusSharing.Groups = append(nodeGpusSharing.Groups, gpuIdx) + nodeGpusSharing.FractionalGpuGroups = append( + nodeGpusSharing.FractionalGpuGroups, + schedulingv1alpha2.FractionalGpuGroup{ + ID: gpuIdx, + ComputeSharingMode: pod.RequestedGPUComputeSharingMode(), + }, + ) } - if len(nodeGpusSharing.Groups) == int(deviceCounts) { + if len(nodeGpusSharing.FractionalGpuGroups) == int(deviceCounts) { return nodeGpusSharing } } @@ -74,7 +93,15 @@ func findGpuForSharingOnNode(task *pod_info.PodInfo, node *node_info.NodeInfo, i isReleasing = false } } - return &nodeGpuForSharing{Groups: []string{string(uuid.NewUUID())}, IsReleasing: isReleasing} + return &nodeGpuForSharing{ + FractionalGpuGroups: []schedulingv1alpha2.FractionalGpuGroup{ + { + ID: string(uuid.NewUUID()), + ComputeSharingMode: task.RequestedGPUComputeSharingMode(), + }, + }, + IsReleasing: isReleasing, + } } func allocateSharedGPUTask(ssn *framework.Session, stmt *framework.Statement, node *node_info.NodeInfo, @@ -83,7 +110,7 @@ func allocateSharedGPUTask(ssn *framework.Session, stmt *framework.Statement, no log.InfraLogger.V(6).Infof( "Pipelining Task <%v/%v> to node <%v> gpuGroup: <%v>, requires: <%v, %v mb> GPUs", task.Namespace, task.Name, node.Name, - task.GPUGroups, task.GpuRequirement.GPUs(), task.GpuRequirement.GpuMemory()) + task.GPUGroupIDs(), task.GpuRequirement.GPUs(), task.GpuRequirement.GpuMemory()) if err := stmt.Pipeline(task, node.Name, !isPipelineOnly); err != nil { log.InfraLogger.V(6).Infof("Failed to pipeline Task: <%s/%s> on Node: <%s>, due to an error: %v", task.Namespace, task.Name, node.Name, err) diff --git a/pkg/scheduler/gpu_sharing/gpuSharing_test.go b/pkg/scheduler/gpu_sharing/gpuSharing_test.go index 51ca324f6..b47a0a67f 100644 --- a/pkg/scheduler/gpu_sharing/gpuSharing_test.go +++ b/pkg/scheduler/gpu_sharing/gpuSharing_test.go @@ -284,15 +284,19 @@ func Test_getNodePreferableGpuForSharing(t *testing.T) { t.Errorf("getNodePreferableGpuForSharing().IsReleasing = %v, want %v", tt.want.isReleasing, gpusForSharing.IsReleasing) } - if len(gpusForSharing.Groups) != tt.want.groupLength { + gpuGroupIDs := make([]string, 0, len(gpusForSharing.FractionalGpuGroups)) + for _, fractionalGpuGroup := range gpusForSharing.FractionalGpuGroups { + gpuGroupIDs = append(gpuGroupIDs, fractionalGpuGroup.ID) + } + if len(gpuGroupIDs) != tt.want.groupLength { t.Errorf("getNodePreferableGpuForSharing() groups array %v, wanted length %v", - gpusForSharing.Groups, tt.want.groupLength) + gpuGroupIDs, tt.want.groupLength) } if tt.want.expectedGroupsInList != nil { for _, expectedGroup := range tt.want.expectedGroupsInList { - if !slices.Contains(gpusForSharing.Groups, expectedGroup) { + if !slices.Contains(gpuGroupIDs, expectedGroup) { t.Errorf("getNodePreferableGpuForSharing() groups array %v, expected to include %v", - gpusForSharing.Groups, expectedGroup) + gpuGroupIDs, expectedGroup) } } } diff --git a/pkg/scheduler/plugins/gpusharingnodevalidation/gpusharingnodevalidation.go b/pkg/scheduler/plugins/gpusharingnodevalidation/gpusharingnodevalidation.go index 892de90a9..cf267a47a 100644 --- a/pkg/scheduler/plugins/gpusharingnodevalidation/gpusharingnodevalidation.go +++ b/pkg/scheduler/plugins/gpusharingnodevalidation/gpusharingnodevalidation.go @@ -108,13 +108,13 @@ func willCreateNewGpuGroup(task *pod_info.PodInfo, node *node_info.NodeInfo, ssn gpuForSharingImmediate := gpu_sharing.GetNodePreferableGpuForSharing(fittingGPUs, node, task, false) if gpuForSharingImmediate != nil && !gpuForSharingImmediate.IsReleasing { - return containsNewGpuGroup(gpuForSharingImmediate.Groups) + return containsNewGpuGroup(gpuForSharingImmediate.GroupIDs()) } gpuForSharingPipelined := gpu_sharing.GetNodePreferableGpuForSharing(fittingGPUs, node, task, true) if gpuForSharingPipelined != nil { - return containsNewGpuGroup(gpuForSharingPipelined.Groups) + return containsNewGpuGroup(gpuForSharingPipelined.GroupIDs()) } // No GPU assignment possible - conservatively assume new group would be needed diff --git a/pkg/scheduler/test_utils/jobs_fake/jobs.go b/pkg/scheduler/test_utils/jobs_fake/jobs.go index d7ff36354..923654870 100644 --- a/pkg/scheduler/test_utils/jobs_fake/jobs.go +++ b/pkg/scheduler/test_utils/jobs_fake/jobs.go @@ -226,7 +226,7 @@ func generateTasks( taskInfo := pod_info.NewTaskInfo(podOfTask, vectorMap, pod_info.TaskInfoOptions{DraPodClaims: draPodClaims}) taskInfo.Status = task.State - taskInfo.GPUGroups = gpuGroups + taskInfo.SetGPUGroupIDs(gpuGroups) taskInfo.SubGroupName = task.SubGroupName taskInfo.IsLegacyMIGtask = task.IsLegacyMigTask taskInfos = append(taskInfos, taskInfo) @@ -243,7 +243,7 @@ func generateTasks( } if tasks_fake.IsTaskStartedStatus(taskInfo.Status) { - gpuName := taskInfo.NodeName + fmt.Sprint(taskInfo.GPUGroups) + gpuName := taskInfo.NodeName + fmt.Sprint(taskInfo.GPUGroupIDs()) if _, ok := allocatedGPUs[gpuName]; !ok { var void interface{} allocatedGPUs[gpuName] = void diff --git a/pkg/scheduler/test_utils/tasks_fake/tasks.go b/pkg/scheduler/test_utils/tasks_fake/tasks.go index 8cd61ae54..fdb4d07ed 100644 --- a/pkg/scheduler/test_utils/tasks_fake/tasks.go +++ b/pkg/scheduler/test_utils/tasks_fake/tasks.go @@ -30,6 +30,7 @@ type TestTaskBasic struct { State pod_status.PodStatus NodeName string // Relevant if job is running NodeAffinityNames []string + Annotations map[string]string PodAffinityLabels map[string]string PodAffinityTopologyKey string PodAntiAffinityTopologyKey string @@ -39,7 +40,6 @@ type TestTaskBasic struct { ResourceClaimTemplates map[string]string ResourceClaimNames []string PersistentVolumeClaimNames []string - Annotations map[string]string } func BuildPod( @@ -87,6 +87,7 @@ func BuildPod( SchedulerName: "kai-scheduler", }, } + maps.Copy(pod.Annotations, task.Annotations) if gpuMemoryMiB > 0 { pod.Annotations[resources.CalcGpuFractionAnnotationForContainer("main")] = resources.GpuMemoryAnnotationToNvFractionsMemoryRequest(gpuMemoryMiB).String() } diff --git a/pkg/scheduler/test_utils/test_utils.go b/pkg/scheduler/test_utils/test_utils.go index 205b6babf..ac8b2e767 100644 --- a/pkg/scheduler/test_utils/test_utils.go +++ b/pkg/scheduler/test_utils/test_utils.go @@ -168,10 +168,11 @@ func MatchExpectedAndRealTasks(t *testing.T, testNumber int, testMetadata TestTo sumOfAcceptedGpus += taskInfo.AcceptedGpuRequirement.GPUs() // verify fractional GPUs index + taskGPUGroups := taskInfo.GPUGroupIDs() if pod_status.IsActiveUsedStatus(taskInfo.Status) && !jobExpectedResult.DontValidateGPUGroup && taskInfo.IsSharedGPUAllocation() && - slices.Equal(taskInfo.GPUGroups, jobExpectedResult.GPUGroups) { + slices.Equal(taskGPUGroups, jobExpectedResult.GPUGroups) { nodeGPUs, found := tasksToGPUGroup[taskInfo.NodeName] if !found { tasksToGPUGroup[taskInfo.NodeName] = make(map[string]string) @@ -179,12 +180,12 @@ func MatchExpectedAndRealTasks(t *testing.T, testNumber int, testMetadata TestTo } for gpuGroupIndex, expectedGpuGroup := range jobExpectedResult.GPUGroups { if gpuGroup, found := nodeGPUs[expectedGpuGroup]; !found { - nodeGPUs[expectedGpuGroup] = taskInfo.GPUGroups[gpuGroupIndex] - } else if gpuGroup != taskInfo.GPUGroups[gpuGroupIndex] { + nodeGPUs[expectedGpuGroup] = taskGPUGroups[gpuGroupIndex] + } else if gpuGroup != taskGPUGroups[gpuGroupIndex] { t.Errorf( "Test number: %d, name: %v, has failed. Task name: %v, "+ "running on GPU: %s, was expecting GPU index: %s", - testNumber, testMetadata.Name, taskInfo.Name, taskInfo.GPUGroups, jobExpectedResult.GPUGroups, + testNumber, testMetadata.Name, taskInfo.Name, taskGPUGroups, jobExpectedResult.GPUGroups, ) } } @@ -291,10 +292,11 @@ func MatchExpectedAndRealTasks(t *testing.T, testNumber int, testMetadata TestTo } // verify fractional GPUs index + taskGPUGroups := task.GPUGroupIDs() if pod_status.IsActiveUsedStatus(task.Status) && !taskExpectedResult.DontValidateGPUGroup && task.IsSharedGPUAllocation() && - slices.Equal(task.GPUGroups, taskExpectedResult.GPUGroups) { + slices.Equal(taskGPUGroups, taskExpectedResult.GPUGroups) { nodeGPUs, found := tasksToGPUGroup[task.NodeName] if !found { tasksToGPUGroup[task.NodeName] = make(map[string]string) @@ -302,12 +304,12 @@ func MatchExpectedAndRealTasks(t *testing.T, testNumber int, testMetadata TestTo } for gpuGroupIndex, expectedGpuGroup := range taskExpectedResult.GPUGroups { if gpuGroup, found := nodeGPUs[expectedGpuGroup]; !found { - nodeGPUs[expectedGpuGroup] = task.GPUGroups[gpuGroupIndex] - } else if gpuGroup != task.GPUGroups[gpuGroupIndex] { + nodeGPUs[expectedGpuGroup] = taskGPUGroups[gpuGroupIndex] + } else if gpuGroup != taskGPUGroups[gpuGroupIndex] { t.Errorf( "Test number: %d, name: %v, has failed. Task name: %v, "+ "running on GPU: %s, was expecting GPU index: %s", - testNumber, testMetadata.Name, taskId, task.GPUGroups, taskExpectedResult.GPUGroups, + testNumber, testMetadata.Name, taskId, taskGPUGroups, taskExpectedResult.GPUGroups, ) } }