Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .changes/unreleased/Fixed-20260728-181225.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
kind: Fixed
body: stalegangeviction no longer evicts remaining pods of a gang whose pods complete successfully
time: 2026-07-28T18:12:25.0211683+05:30
custom:
Author: Thezone-1
Issue: "1968"
195 changes: 195 additions & 0 deletions pkg/scheduler/actions/stalegangeviction/stalegangeviction_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -709,6 +709,201 @@ func TestStaleGangEviction(t *testing.T) {
},
},
},
{
name: "Gang with a succeeded pod - no evict",
topology: test_utils.TestTopologyBasic{
Jobs: []*jobs_fake.TestJobBasic{
{
Name: "job-1",
QueueName: "q-1",
RootSubGroupSet: jobs_fake.DefaultSubGroup(3),
Tasks: []*tasks_fake.TestTaskBasic{
{
Name: "job-1-0",
State: pod_status.Succeeded,
},
{
Name: "job-1-1",
State: pod_status.Running,
NodeName: "node-1",
},
{
Name: "job-1-2",
State: pod_status.Running,
NodeName: "node-1",
},
},
StaleDuration: pointer.Duration(61 * time.Second),
},
},
Nodes: map[string]nodes_fake.TestNodeBasic{
"node-1": {},
},
Queues: []test_utils.TestQueueBasic{
{
Name: "q-1",
ParentQueue: "d-1",
},
},
Departments: []test_utils.TestDepartmentBasic{
{
Name: "d-1",
},
},
TaskExpectedResults: map[string]test_utils.TestExpectedResultBasic{
"job-1-0": {
Status: pod_status.Succeeded,
},
"job-1-1": {
NodeName: "node-1",
Status: pod_status.Running,
},
"job-1-2": {
NodeName: "node-1",
Status: pod_status.Running,
},
},
Mocks: &test_utils.TestMock{
CacheRequirements: &test_utils.CacheMocking{
NumberOfCacheBinds: 0,
NumberOfCacheEvictions: 0,
NumberOfPipelineActions: 0,
},
},
},
},
{
name: "Gang with only one pod left running - no evict",
topology: test_utils.TestTopologyBasic{
Jobs: []*jobs_fake.TestJobBasic{
{
Name: "job-1",
QueueName: "q-1",
RootSubGroupSet: jobs_fake.DefaultSubGroup(3),
Tasks: []*tasks_fake.TestTaskBasic{
{
Name: "job-1-0",
State: pod_status.Succeeded,
},
{
Name: "job-1-1",
State: pod_status.Succeeded,
},
{
Name: "job-1-2",
State: pod_status.Running,
NodeName: "node-1",
},
},
StaleDuration: pointer.Duration(61 * time.Second),
},
},
Nodes: map[string]nodes_fake.TestNodeBasic{
"node-1": {},
},
Queues: []test_utils.TestQueueBasic{
{
Name: "q-1",
ParentQueue: "d-1",
},
},
Departments: []test_utils.TestDepartmentBasic{
{
Name: "d-1",
},
},
TaskExpectedResults: map[string]test_utils.TestExpectedResultBasic{
"job-1-0": {
Status: pod_status.Succeeded,
},
"job-1-1": {
Status: pod_status.Succeeded,
},
"job-1-2": {
NodeName: "node-1",
Status: pod_status.Running,
},
},
Mocks: &test_utils.TestMock{
CacheRequirements: &test_utils.CacheMocking{
NumberOfCacheBinds: 0,
NumberOfCacheEvictions: 0,
NumberOfPipelineActions: 0,
},
},
},
},
{
name: "Gang with a succeeded pod in one sub group - no evict",
topology: test_utils.TestTopologyBasic{
Jobs: []*jobs_fake.TestJobBasic{
{
Name: "job-1",
QueueName: "q-1",
RootSubGroupSet: func() *subgroup_info.SubGroupSet {
root := subgroup_info.NewSubGroupSet(subgroup_info.RootSubGroupSetName, nil)
root.AddPodSet(subgroup_info.NewPodSet("sub-group-0", 2, nil))
root.AddPodSet(subgroup_info.NewPodSet("sub-group-1", 1, nil))
return root
}(),
Tasks: []*tasks_fake.TestTaskBasic{
{
Name: "job-1-0",
SubGroupName: "sub-group-0",
State: pod_status.Succeeded,
},
{
Name: "job-1-1",
SubGroupName: "sub-group-0",
State: pod_status.Running,
NodeName: "node-1",
},
{
Name: "job-1-2",
SubGroupName: "sub-group-1",
State: pod_status.Running,
NodeName: "node-1",
},
},
StaleDuration: pointer.Duration(61 * time.Second),
},
},
Nodes: map[string]nodes_fake.TestNodeBasic{
"node-1": {},
},
Queues: []test_utils.TestQueueBasic{
{
Name: "q-1",
ParentQueue: "d-1",
},
},
Departments: []test_utils.TestDepartmentBasic{
{
Name: "d-1",
},
},
TaskExpectedResults: map[string]test_utils.TestExpectedResultBasic{
"job-1-0": {
Status: pod_status.Succeeded,
},
"job-1-1": {
NodeName: "node-1",
Status: pod_status.Running,
},
"job-1-2": {
NodeName: "node-1",
Status: pod_status.Running,
},
},
Mocks: &test_utils.TestMock{
CacheRequirements: &test_utils.CacheMocking{
NumberOfCacheBinds: 0,
NumberOfCacheEvictions: 0,
NumberOfPipelineActions: 0,
},
},
},
},
} {
t.Run(test.name, func(t *testing.T) {
t.Logf("Running test number: %v, test name: %v,", i, test.name)
Expand Down
12 changes: 8 additions & 4 deletions pkg/scheduler/cache/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,9 +83,13 @@ var terminalPodPhases = []v1.PodPhase{
v1.PodFailed,
}

func filterTerminalPods(options *metav1.ListOptions) {
selectors := make([]string, 0, len(terminalPodPhases))
for _, phase := range terminalPodPhases {
var watchFilteredPodPhases = []v1.PodPhase{
v1.PodFailed,
}

func filterFailedPods(options *metav1.ListOptions) {
selectors := make([]string, 0, len(watchFilteredPodPhases))
for _, phase := range watchFilteredPodPhases {
selectors = append(selectors, fmt.Sprintf("status.phase!=%s", phase))
}
selector := strings.Join(selectors, ",")
Expand All @@ -103,7 +107,7 @@ func registerSchedulerPodInformer(informerFactory informers.SharedInformerFactor
metav1.NamespaceAll,
resyncPeriod,
k8scache.Indexers{k8scache.NamespaceIndex: k8scache.MetaNamespaceIndexFunc},
filterTerminalPods,
filterFailedPods,
)
})
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/scheduler/cache/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ func TestCache(t *testing.T) {
var _ = Describe("Cache", func() {
Describe("New", func() {
Context("Pod informer filtering", func() {
It("should filter terminal pods without filtering pods by scheduler name", func() {
It("should filter failed pods while keeping succeeded pods, without filtering by scheduler name", func() {
kubeClient := fake.NewSimpleClientset()
cache := New(&SchedulerCacheParams{
KubeClient: kubeClient,
Expand Down Expand Up @@ -95,8 +95,8 @@ var _ = Describe("Cache", func() {

Expect(podSelectors).NotTo(BeEmpty())
for _, selector := range podSelectors {
Expect(selector).To(ContainSubstring("status.phase!=Succeeded"))
Expect(selector).To(ContainSubstring("status.phase!=Failed"))
Expect(selector).NotTo(ContainSubstring("status.phase!=Succeeded"))
Expect(selector).NotTo(ContainSubstring("spec.schedulerName"))
}
Expect(nonPodSelectors).To(BeEmpty())
Expand Down
29 changes: 29 additions & 0 deletions pkg/scheduler/cache/pod_transform.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,10 @@ package cache

import (
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/tools/cache"

commonconstants "github.com/kai-scheduler/KAI-scheduler/pkg/common/constants"
)

func setSchedulerPodTransform(informer cache.SharedIndexInformer) error {
Expand All @@ -18,6 +21,10 @@ func compactSchedulerPod(obj any) (any, error) {
return obj, nil
}

if pod.Status.Phase == v1.PodSucceeded {
return compactSucceededPod(pod), nil
}

compact := pod.DeepCopy()
compact.ManagedFields = nil
compact.Spec.Containers = compactContainers(compact.Spec.Containers)
Expand All @@ -26,6 +33,28 @@ func compactSchedulerPod(obj any) (any, error) {
return compact, nil
}

func compactSucceededPod(pod *v1.Pod) *v1.Pod {
compact := &v1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: pod.Name,
Namespace: pod.Namespace,
UID: pod.UID,
},
Status: v1.PodStatus{
Phase: pod.Status.Phase,
},
}

if podGroup, found := pod.Annotations[commonconstants.PodGroupAnnotationForPod]; found {
compact.Annotations = map[string]string{commonconstants.PodGroupAnnotationForPod: podGroup}
}
if subGroup, found := pod.Labels[commonconstants.SubGroupLabelKey]; found {
compact.Labels = map[string]string{commonconstants.SubGroupLabelKey: subGroup}
}

return compact
}

func compactContainers(containers []v1.Container) []v1.Container {
compact := make([]v1.Container, 0, len(containers))
for _, container := range containers {
Expand Down
Loading
Loading