Move GetAffinityTerms functions from pkg/scheduler/framework to staging repo

This commit is contained in:
Ania Borowiec
2025-08-26 13:05:17 +00:00
parent 091f87c10b
commit 3c00c3cb29
5 changed files with 153 additions and 145 deletions

View File

@@ -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
}

View File

@@ -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.

View File

@@ -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

View File

@@ -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

View File

@@ -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)
}
})
}
}