From 3c00c3cb29d1ecbed9221d268282744bf183e517 Mon Sep 17 00:00:00 2001 From: Ania Borowiec Date: Tue, 26 Aug 2025 13:05:17 +0000 Subject: [PATCH] Move GetAffinityTerms functions from pkg/scheduler/framework to staging repo --- .../plugins/interpodaffinity/plugin.go | 8 +- pkg/scheduler/framework/types.go | 101 +----------------- pkg/scheduler/framework/types_test.go | 43 -------- .../k8s.io/kube-scheduler/framework/types.go | 94 +++++++++++++++- .../kube-scheduler/framework/types_test.go | 52 ++++++++- 5 files changed, 153 insertions(+), 145 deletions(-) diff --git a/pkg/scheduler/framework/plugins/interpodaffinity/plugin.go b/pkg/scheduler/framework/plugins/interpodaffinity/plugin.go index f3f6c89e3d5..f2f66bd23ab 100644 --- a/pkg/scheduler/framework/plugins/interpodaffinity/plugin.go +++ b/pkg/scheduler/framework/plugins/interpodaffinity/plugin.go @@ -159,12 +159,12 @@ func (pl *InterPodAffinity) isSchedulableAfterPodChange(logger klog.Logger, pod return fwk.QueueSkip, nil } - terms, err := framework.GetAffinityTerms(pod, framework.GetPodAffinityTerms(pod.Spec.Affinity)) + terms, err := fwk.GetAffinityTerms(pod, fwk.GetPodAffinityTerms(pod.Spec.Affinity)) if err != nil { return fwk.Queue, err } - antiTerms, err := framework.GetAffinityTerms(pod, framework.GetPodAntiAffinityTerms(pod.Spec.Affinity)) + antiTerms, err := fwk.GetAffinityTerms(pod, fwk.GetPodAntiAffinityTerms(pod.Spec.Affinity)) if err != nil { return fwk.Queue, err } @@ -217,7 +217,7 @@ func (pl *InterPodAffinity) isSchedulableAfterNodeChange(logger klog.Logger, pod return fwk.Queue, err } - terms, err := framework.GetAffinityTerms(pod, framework.GetPodAffinityTerms(pod.Spec.Affinity)) + terms, err := fwk.GetAffinityTerms(pod, fwk.GetPodAffinityTerms(pod.Spec.Affinity)) if err != nil { return fwk.Queue, err } @@ -254,7 +254,7 @@ func (pl *InterPodAffinity) isSchedulableAfterNodeChange(logger klog.Logger, pod } } - antiTerms, err := framework.GetAffinityTerms(pod, framework.GetPodAntiAffinityTerms(pod.Spec.Affinity)) + antiTerms, err := fwk.GetAffinityTerms(pod, fwk.GetPodAntiAffinityTerms(pod.Spec.Affinity)) if err != nil { return fwk.Queue, err } diff --git a/pkg/scheduler/framework/types.go b/pkg/scheduler/framework/types.go index 969b3898f83..f24f516e720 100644 --- a/pkg/scheduler/framework/types.go +++ b/pkg/scheduler/framework/types.go @@ -26,7 +26,6 @@ import ( "time" v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" utilerrors "k8s.io/apimachinery/pkg/util/errors" utilruntime "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/apimachinery/pkg/util/sets" @@ -670,20 +669,20 @@ func (pi *PodInfo) Update(pod *v1.Pod) error { // Attempt to parse the affinity terms var parseErrs []error - requiredAffinityTerms, err := GetAffinityTerms(pod, GetPodAffinityTerms(pod.Spec.Affinity)) + requiredAffinityTerms, err := fwk.GetAffinityTerms(pod, fwk.GetPodAffinityTerms(pod.Spec.Affinity)) if err != nil { parseErrs = append(parseErrs, fmt.Errorf("requiredAffinityTerms: %w", err)) } - requiredAntiAffinityTerms, err := GetAffinityTerms(pod, - GetPodAntiAffinityTerms(pod.Spec.Affinity)) + requiredAntiAffinityTerms, err := fwk.GetAffinityTerms(pod, + fwk.GetPodAntiAffinityTerms(pod.Spec.Affinity)) if err != nil { parseErrs = append(parseErrs, fmt.Errorf("requiredAntiAffinityTerms: %w", err)) } - weightedAffinityTerms, err := getWeightedAffinityTerms(pod, preferredAffinityTerms) + weightedAffinityTerms, err := fwk.GetWeightedAffinityTerms(pod, preferredAffinityTerms) if err != nil { parseErrs = append(parseErrs, fmt.Errorf("preferredAffinityTerms: %w", err)) } - weightedAntiAffinityTerms, err := getWeightedAffinityTerms(pod, preferredAntiAffinityTerms) + weightedAntiAffinityTerms, err := fwk.GetWeightedAffinityTerms(pod, preferredAntiAffinityTerms) if err != nil { parseErrs = append(parseErrs, fmt.Errorf("preferredAntiAffinityTerms: %w", err)) } @@ -837,58 +836,6 @@ func (f *FitError) Error() string { return reasonMsg } -func newAffinityTerm(pod *v1.Pod, term *v1.PodAffinityTerm) (*fwk.AffinityTerm, error) { - selector, err := metav1.LabelSelectorAsSelector(term.LabelSelector) - if err != nil { - return nil, err - } - - namespaces := getNamespacesFromPodAffinityTerm(pod, term) - nsSelector, err := metav1.LabelSelectorAsSelector(term.NamespaceSelector) - if err != nil { - return nil, err - } - - return &fwk.AffinityTerm{Namespaces: namespaces, Selector: selector, TopologyKey: term.TopologyKey, NamespaceSelector: nsSelector}, nil -} - -// GetAffinityTerms receives a Pod and affinity terms and returns the namespaces and -// selectors of the terms. -func GetAffinityTerms(pod *v1.Pod, v1Terms []v1.PodAffinityTerm) ([]fwk.AffinityTerm, error) { - if v1Terms == nil { - return nil, nil - } - - var terms []fwk.AffinityTerm - for i := range v1Terms { - t, err := newAffinityTerm(pod, &v1Terms[i]) - if err != nil { - // We get here if the label selector failed to process - return nil, err - } - terms = append(terms, *t) - } - return terms, nil -} - -// getWeightedAffinityTerms returns the list of processed affinity terms. -func getWeightedAffinityTerms(pod *v1.Pod, v1Terms []v1.WeightedPodAffinityTerm) ([]fwk.WeightedAffinityTerm, error) { - if v1Terms == nil { - return nil, nil - } - - var terms []fwk.WeightedAffinityTerm - for i := range v1Terms { - t, err := newAffinityTerm(pod, &v1Terms[i].PodAffinityTerm) - if err != nil { - // We get here if the label selector failed to process - return nil, err - } - terms = append(terms, fwk.WeightedAffinityTerm{AffinityTerm: *t, Weight: v1Terms[i].Weight}) - } - return terms, nil -} - // NewPodInfo returns a new PodInfo. func NewPodInfo(pod *v1.Pod) (*PodInfo, error) { pInfo := &PodInfo{} @@ -896,44 +843,6 @@ func NewPodInfo(pod *v1.Pod) (*PodInfo, error) { return pInfo, err } -func GetPodAffinityTerms(affinity *v1.Affinity) (terms []v1.PodAffinityTerm) { - if affinity != nil && affinity.PodAffinity != nil { - if len(affinity.PodAffinity.RequiredDuringSchedulingIgnoredDuringExecution) != 0 { - terms = affinity.PodAffinity.RequiredDuringSchedulingIgnoredDuringExecution - } - // TODO: Uncomment this block when implement RequiredDuringSchedulingRequiredDuringExecution. - // if len(affinity.PodAffinity.RequiredDuringSchedulingRequiredDuringExecution) != 0 { - // terms = append(terms, affinity.PodAffinity.RequiredDuringSchedulingRequiredDuringExecution...) - // } - } - return terms -} - -func GetPodAntiAffinityTerms(affinity *v1.Affinity) (terms []v1.PodAffinityTerm) { - if affinity != nil && affinity.PodAntiAffinity != nil { - if len(affinity.PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution) != 0 { - terms = affinity.PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution - } - // TODO: Uncomment this block when implement RequiredDuringSchedulingRequiredDuringExecution. - // if len(affinity.PodAntiAffinity.RequiredDuringSchedulingRequiredDuringExecution) != 0 { - // terms = append(terms, affinity.PodAntiAffinity.RequiredDuringSchedulingRequiredDuringExecution...) - // } - } - return terms -} - -// returns a set of names according to the namespaces indicated in podAffinityTerm. -// If namespaces is empty it considers the given pod's namespace. -func getNamespacesFromPodAffinityTerm(pod *v1.Pod, podAffinityTerm *v1.PodAffinityTerm) sets.Set[string] { - names := sets.Set[string]{} - if len(podAffinityTerm.Namespaces) == 0 && podAffinityTerm.NamespaceSelector == nil { - names.Insert(pod.Namespace) - } else { - names.Insert(podAffinityTerm.Namespaces...) - } - return names -} - // Resource is a collection of compute resource. // Implementation is separate from interface fwk.Resource, because implementation of functions Add and SetMaxResource // depends on internal scheduler util functions. diff --git a/pkg/scheduler/framework/types_test.go b/pkg/scheduler/framework/types_test.go index 5b6bc19ecb9..15b09d57526 100644 --- a/pkg/scheduler/framework/types_test.go +++ b/pkg/scheduler/framework/types_test.go @@ -27,7 +27,6 @@ import ( "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" - "k8s.io/apimachinery/pkg/util/sets" utilfeature "k8s.io/apiserver/pkg/util/feature" featuregatetesting "k8s.io/component-base/featuregate/testing" "k8s.io/klog/v2" @@ -1227,48 +1226,6 @@ func fakeNodeInfo(pods ...*v1.Pod) *NodeInfo { return ni } -func TestGetNamespacesFromPodAffinityTerm(t *testing.T) { - tests := []struct { - name string - term *v1.PodAffinityTerm - want sets.Set[string] - }{ - { - name: "podAffinityTerm_namespace_empty", - term: &v1.PodAffinityTerm{}, - want: sets.Set[string]{metav1.NamespaceDefault: sets.Empty{}}, - }, - { - name: "podAffinityTerm_namespace_not_empty", - term: &v1.PodAffinityTerm{ - Namespaces: []string{metav1.NamespacePublic, metav1.NamespaceSystem}, - }, - want: sets.New(metav1.NamespacePublic, metav1.NamespaceSystem), - }, - { - name: "podAffinityTerm_namespace_selector_not_nil", - term: &v1.PodAffinityTerm{ - NamespaceSelector: &metav1.LabelSelector{}, - }, - want: sets.Set[string]{}, - }, - } - - for _, test := range tests { - t.Run(test.name, func(t *testing.T) { - got := getNamespacesFromPodAffinityTerm(&v1.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: "topologies_pod", - Namespace: metav1.NamespaceDefault, - }, - }, test.term) - if diff := cmp.Diff(test.want, got); diff != "" { - t.Errorf("Unexpected namespaces (-want, +got):\n%s", diff) - } - }) - } -} - func TestFitError_Error(t *testing.T) { tests := []struct { name string diff --git a/staging/src/k8s.io/kube-scheduler/framework/types.go b/staging/src/k8s.io/kube-scheduler/framework/types.go index 8b9a7e1dfbb..135d00a32bc 100644 --- a/staging/src/k8s.io/kube-scheduler/framework/types.go +++ b/staging/src/k8s.io/kube-scheduler/framework/types.go @@ -21,9 +21,9 @@ import ( "time" v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/util/sets" - "k8s.io/klog/v2" ) @@ -389,6 +389,98 @@ type WeightedAffinityTerm struct { Weight int32 } +// GetAffinityTerms receives a Pod and affinity terms and returns the namespaces and +// selectors of the terms. +func GetAffinityTerms(pod *v1.Pod, v1Terms []v1.PodAffinityTerm) ([]AffinityTerm, error) { + if v1Terms == nil { + return nil, nil + } + + var terms []AffinityTerm + for i := range v1Terms { + t, err := newAffinityTerm(pod, &v1Terms[i]) + if err != nil { + // We get here if the label selector failed to process + return nil, err + } + terms = append(terms, *t) + } + return terms, nil +} + +func newAffinityTerm(pod *v1.Pod, term *v1.PodAffinityTerm) (*AffinityTerm, error) { + selector, err := metav1.LabelSelectorAsSelector(term.LabelSelector) + if err != nil { + return nil, err + } + + namespaces := getNamespacesFromPodAffinityTerm(pod, term) + nsSelector, err := metav1.LabelSelectorAsSelector(term.NamespaceSelector) + if err != nil { + return nil, err + } + + return &AffinityTerm{Namespaces: namespaces, Selector: selector, TopologyKey: term.TopologyKey, NamespaceSelector: nsSelector}, nil +} + +// returns a set of names according to the namespaces indicated in podAffinityTerm. +// If namespaces is empty it considers the given pod's namespace. +func getNamespacesFromPodAffinityTerm(pod *v1.Pod, podAffinityTerm *v1.PodAffinityTerm) sets.Set[string] { + names := sets.Set[string]{} + if len(podAffinityTerm.Namespaces) == 0 && podAffinityTerm.NamespaceSelector == nil { + names.Insert(pod.Namespace) + } else { + names.Insert(podAffinityTerm.Namespaces...) + } + return names +} + +// Returns the list of PodAffinityTerms specified in the PodAffinity.RequiredDuringSchedulingIgnoredDuringExecution field. +func GetPodAffinityTerms(affinity *v1.Affinity) (terms []v1.PodAffinityTerm) { + if affinity != nil && affinity.PodAffinity != nil { + if len(affinity.PodAffinity.RequiredDuringSchedulingIgnoredDuringExecution) != 0 { + terms = affinity.PodAffinity.RequiredDuringSchedulingIgnoredDuringExecution + } + // TODO: Uncomment this block when implement RequiredDuringSchedulingRequiredDuringExecution. + // if len(affinity.PodAffinity.RequiredDuringSchedulingRequiredDuringExecution) != 0 { + // terms = append(terms, affinity.PodAffinity.RequiredDuringSchedulingRequiredDuringExecution...) + // } + } + return terms +} + +// GetWeightedAffinityTerms returns the list of processed affinity terms. +func GetWeightedAffinityTerms(pod *v1.Pod, v1Terms []v1.WeightedPodAffinityTerm) ([]WeightedAffinityTerm, error) { + if v1Terms == nil { + return nil, nil + } + + var terms []WeightedAffinityTerm + for i := range v1Terms { + t, err := newAffinityTerm(pod, &v1Terms[i].PodAffinityTerm) + if err != nil { + // We get here if the label selector failed to process + return nil, err + } + terms = append(terms, WeightedAffinityTerm{AffinityTerm: *t, Weight: v1Terms[i].Weight}) + } + return terms, nil +} + +// Returns the list of PodAffinityTerms specified in the PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution field. +func GetPodAntiAffinityTerms(affinity *v1.Affinity) (terms []v1.PodAffinityTerm) { + if affinity != nil && affinity.PodAntiAffinity != nil { + if len(affinity.PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution) != 0 { + terms = affinity.PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution + } + // TODO: Uncomment this block when implement RequiredDuringSchedulingRequiredDuringExecution. + // if len(affinity.PodAntiAffinity.RequiredDuringSchedulingRequiredDuringExecution) != 0 { + // terms = append(terms, affinity.PodAntiAffinity.RequiredDuringSchedulingRequiredDuringExecution...) + // } + } + return terms +} + // Resource is a collection of compute resources. type Resource interface { GetMilliCPU() int64 diff --git a/staging/src/k8s.io/kube-scheduler/framework/types_test.go b/staging/src/k8s.io/kube-scheduler/framework/types_test.go index d9ae25340fe..8f4bfc3f4ac 100644 --- a/staging/src/k8s.io/kube-scheduler/framework/types_test.go +++ b/staging/src/k8s.io/kube-scheduler/framework/types_test.go @@ -16,7 +16,15 @@ limitations under the License. package framework -import "testing" +import ( + "testing" + + "github.com/google/go-cmp/cmp" + + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/sets" +) type hostPortInfoParam struct { protocol, ip string @@ -231,3 +239,45 @@ func TestHostPortInfo_Check(t *testing.T) { }) } } + +func TestGetNamespacesFromPodAffinityTerm(t *testing.T) { + tests := []struct { + name string + term *v1.PodAffinityTerm + want sets.Set[string] + }{ + { + name: "podAffinityTerm_namespace_empty", + term: &v1.PodAffinityTerm{}, + want: sets.Set[string]{metav1.NamespaceDefault: sets.Empty{}}, + }, + { + name: "podAffinityTerm_namespace_not_empty", + term: &v1.PodAffinityTerm{ + Namespaces: []string{metav1.NamespacePublic, metav1.NamespaceSystem}, + }, + want: sets.New(metav1.NamespacePublic, metav1.NamespaceSystem), + }, + { + name: "podAffinityTerm_namespace_selector_not_nil", + term: &v1.PodAffinityTerm{ + NamespaceSelector: &metav1.LabelSelector{}, + }, + want: sets.Set[string]{}, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got := getNamespacesFromPodAffinityTerm(&v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: "topologies_pod", + Namespace: metav1.NamespaceDefault, + }, + }, test.term) + if diff := cmp.Diff(test.want, got); diff != "" { + t.Errorf("Unexpected namespaces (-want, +got):\n%s", diff) + } + }) + } +}