Skip to content

Commit 2f78239

Browse files
authored
fix(scheduler): fail fast when feature gate discovery fails (#2050)
Signed-off-by: Huy Nguyen VN <hcnguyen@nvidia.com>
1 parent 46a09dd commit 2f78239

11 files changed

Lines changed: 176 additions & 44 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
kind: Fixed
2+
body: |-
3+
Scheduler and binder now fail fast when Kubernetes API discovery fails at startup, instead of silently disabling dynamic resource allocation and node resource topology for the lifetime of the process.
4+
custom:
5+
Issue: "2038"
6+
Author: hcnguyen-ai

cmd/binder/app/app.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,7 +105,10 @@ func New(options *Options, config *rest.Config) (*App, error) {
105105
kubeClient := draversionawareclient.NewDRAAwareClient(kubernetes.NewForConfigOrDie(config))
106106
informerFactory := informers.NewSharedInformerFactory(kubeClient, 0)
107107

108-
featuregates.SetDRAFeatureGate(kubeClient.Discovery())
108+
if err := featuregates.SetDRAFeatureGate(kubeClient.Discovery()); err != nil {
109+
setupLog.Error(err, "unable to determine dynamic resource allocation availability")
110+
return nil, err
111+
}
109112

110113
rrs := resourcereservation.NewService(options.FakeGPUNodes, clientWithWatch, options.ResourceReservationPodImage,
111114
time.Duration(options.ResourceReservationAllocationTimeout)*time.Second,

cmd/snapshot-tool/main.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,10 @@ func main() {
9393
DiscoveryClient: kubeClient.Discovery(),
9494
}
9595

96-
schedulerCache := cache.New(schedulerCacheParams)
96+
schedulerCache, err := cache.New(schedulerCacheParams)
97+
if err != nil {
98+
log.InfraLogger.Fatalf("Failed to create scheduler cache: %v", err)
99+
}
97100
stopCh := make(chan struct{})
98101
schedulerCache.Run(stopCh)
99102
schedulerCache.WaitForCacheSync(stopCh)

pkg/common/feature_gates/feature_gates.go

Lines changed: 26 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -4,14 +4,14 @@
44
package featuregates
55

66
import (
7+
"fmt"
78
"strconv"
89
"strings"
910
"sync/atomic"
1011

1112
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1213
"k8s.io/apimachinery/pkg/version"
1314
discovery "k8s.io/client-go/discovery"
14-
"sigs.k8s.io/controller-runtime/pkg/log"
1515
)
1616

1717
const (
@@ -27,8 +27,13 @@ const (
2727
// off to reflect server-side DRA availability.
2828
var dynamicResourcesEnabled atomic.Bool
2929

30-
func SetDRAFeatureGate(discoveryClient discovery.DiscoveryInterface) {
31-
dynamicResourcesEnabled.Store(IsDynamicResourcesEnabled(discoveryClient))
30+
func SetDRAFeatureGate(discoveryClient discovery.DiscoveryInterface) error {
31+
enabled, err := IsDynamicResourcesEnabled(discoveryClient)
32+
if err != nil {
33+
return err
34+
}
35+
dynamicResourcesEnabled.Store(enabled)
36+
return nil
3237
}
3338

3439
// DynamicResourcesEnabled reports whether DRA was determined to be usable
@@ -47,8 +52,13 @@ func SetDynamicResourcesEnabledForTest(enabled bool) {
4752

4853
var nodeResourceTopologyEnabled atomic.Bool
4954

50-
func SetNodeResourceTopologyFeatureGate(discoveryClient discovery.DiscoveryInterface) {
51-
nodeResourceTopologyEnabled.Store(IsNodeResourceTopologyEnabled(discoveryClient))
55+
func SetNodeResourceTopologyFeatureGate(discoveryClient discovery.DiscoveryInterface) error {
56+
enabled, err := IsNodeResourceTopologyEnabled(discoveryClient)
57+
if err != nil {
58+
return err
59+
}
60+
nodeResourceTopologyEnabled.Store(enabled)
61+
return nil
5262
}
5363

5464
func NodeResourceTopologyEnabled() bool {
@@ -61,43 +71,36 @@ func SetNodeResourceTopologyEnabledForTest(enabled bool) {
6171

6272
// IsNodeResourceTopologyEnabled reports whether the cluster serves the
6373
// topology.node.k8s.io API group (the NodeResourceTopology CRD).
64-
func IsNodeResourceTopologyEnabled(discoveryClient discovery.DiscoveryInterface) bool {
65-
logger := log.Log.WithName("feature-gates")
66-
74+
func IsNodeResourceTopologyEnabled(discoveryClient discovery.DiscoveryInterface) (bool, error) {
6775
serverGroups, err := discoveryClient.ServerGroups()
6876
if err != nil {
69-
logger.Error(err, "Failed to get server groups")
70-
return false
77+
return false, fmt.Errorf("failed to get server groups: %w", err)
7178
}
7279

7380
for _, group := range serverGroups.Groups {
7481
if group.Name == nodeResourceTopologyGroup {
75-
return true
82+
return true, nil
7683
}
7784
}
78-
return false
85+
return false, nil
7986
}
8087

81-
func IsDynamicResourcesEnabled(discoveryClient discovery.DiscoveryInterface) bool {
82-
logger := log.Log.WithName("feature-gates")
83-
88+
func IsDynamicResourcesEnabled(discoveryClient discovery.DiscoveryInterface) (bool, error) {
8489
// Get API server version
8590
serverVersion, err := discoveryClient.ServerVersion()
8691
if err != nil {
87-
logger.Error(err, "Failed to get server version")
88-
return false
92+
return false, fmt.Errorf("failed to get server version: %w", err)
8993
}
9094

9195
// Check if the API server version is compatible with DRA
9296
if !isCompatibleDRAVersion(serverVersion) {
93-
return false
97+
return false, nil
9498
}
9599

96100
// Get supported API versions
97101
serverGroups, err := discoveryClient.ServerGroups()
98102
if err != nil {
99-
logger.Error(err, "Failed to get server groups")
100-
return false
103+
return false, fmt.Errorf("failed to get server groups: %w", err)
101104
}
102105

103106
found := false
@@ -110,17 +113,17 @@ func IsDynamicResourcesEnabled(discoveryClient discovery.DiscoveryInterface) boo
110113
}
111114
}
112115
if !found {
113-
return false
116+
return false, nil
114117
}
115118

116119
// Check if the DRA API group is supported
117120
for _, groupVersion := range resourceGroup.Versions {
118121
if version.CompareKubeAwareVersionStrings(groupVersion.Version, minimalSupportedVersion) >= 0 {
119-
return true
122+
return true, nil
120123
}
121124
}
122125

123-
return false
126+
return false, nil
124127
}
125128

126129
func isCompatibleDRAVersion(serverVersion *version.Info) bool {

pkg/common/feature_gates/feature_gates_test.go

Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
package featuregates
55

66
import (
7+
"errors"
78
"testing"
89

910
. "github.com/onsi/ginkgo/v2"
@@ -13,10 +14,42 @@ import (
1314
resourcev1beta1 "k8s.io/api/resource/v1beta1"
1415
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1516
version "k8s.io/apimachinery/pkg/version"
17+
"k8s.io/client-go/discovery"
1618
fakediscovery "k8s.io/client-go/discovery/fake"
1719
"k8s.io/client-go/kubernetes/fake"
1820
)
1921

22+
// erroringDiscovery fails the requested discovery call, standing in for an API
23+
// server that is briefly unreachable while a component starts up.
24+
type erroringDiscovery struct {
25+
discovery.DiscoveryInterface
26+
versionErr error
27+
groupsErr error
28+
}
29+
30+
func (e *erroringDiscovery) ServerVersion() (*version.Info, error) {
31+
if e.versionErr != nil {
32+
return nil, e.versionErr
33+
}
34+
return e.DiscoveryInterface.ServerVersion()
35+
}
36+
37+
func (e *erroringDiscovery) ServerGroups() (*metav1.APIGroupList, error) {
38+
if e.groupsErr != nil {
39+
return nil, e.groupsErr
40+
}
41+
return e.DiscoveryInterface.ServerGroups()
42+
}
43+
44+
func discoveryForVersion(major, minor string) discovery.DiscoveryInterface {
45+
fakeClient := fake.NewClientset()
46+
fakeClient.Discovery().(*fakediscovery.FakeDiscovery).FakedServerVersion = &version.Info{
47+
Major: major,
48+
Minor: minor,
49+
}
50+
return fakeClient.Discovery()
51+
}
52+
2053
func TestCache(t *testing.T) {
2154
RegisterFailHandler(Fail)
2255
RunSpecs(t, "Test cache")
@@ -46,5 +79,75 @@ var _ = Describe("New", func() {
4679
Entry("edge case version (1.31) with resource API should not enable DRA", "1", "31", []string{resourcev1alhpa3.SchemeGroupVersion.String()}, false),
4780
Entry("higher compatible version (1.35) with resource API should enable DRA", "1", "34", []string{resourcev1.SchemeGroupVersion.String()}, true),
4881
)
82+
83+
It("should return an error when the server version cannot be retrieved", func() {
84+
discoveryClient := &erroringDiscovery{
85+
DiscoveryInterface: discoveryForVersion("1", "34"),
86+
versionErr: errors.New("connection refused"),
87+
}
88+
89+
_, err := IsDynamicResourcesEnabled(discoveryClient)
90+
Expect(err).To(HaveOccurred())
91+
})
92+
93+
It("should return an error when the server groups cannot be retrieved", func() {
94+
discoveryClient := &erroringDiscovery{
95+
DiscoveryInterface: discoveryForVersion("1", "34"),
96+
groupsErr: errors.New("connection refused"),
97+
}
98+
99+
_, err := IsDynamicResourcesEnabled(discoveryClient)
100+
Expect(err).To(HaveOccurred())
101+
})
102+
103+
It("should leave the gate untouched when discovery fails", func() {
104+
SetDynamicResourcesEnabledForTest(true)
105+
DeferCleanup(SetDynamicResourcesEnabledForTest, false)
106+
107+
discoveryClient := &erroringDiscovery{
108+
DiscoveryInterface: discoveryForVersion("1", "34"),
109+
versionErr: errors.New("connection refused"),
110+
}
111+
112+
Expect(SetDRAFeatureGate(discoveryClient)).NotTo(Succeed())
113+
Expect(DynamicResourcesEnabled()).To(BeTrue())
114+
})
115+
})
116+
117+
Context("NodeResourceTopology Feature Gate", func() {
118+
It("should report availability when the API group is served", func() {
119+
fakeClient := fake.NewClientset()
120+
fakeClient.Resources = append(fakeClient.Resources,
121+
&metav1.APIResourceList{GroupVersion: nodeResourceTopologyGroup + "/v1alpha2"})
122+
123+
Expect(IsNodeResourceTopologyEnabled(fakeClient.Discovery())).To(BeTrue())
124+
})
125+
126+
It("should report unavailability when the API group is absent", func() {
127+
Expect(IsNodeResourceTopologyEnabled(fake.NewClientset().Discovery())).To(BeFalse())
128+
})
129+
130+
It("should return an error when the server groups cannot be retrieved", func() {
131+
discoveryClient := &erroringDiscovery{
132+
DiscoveryInterface: fake.NewClientset().Discovery(),
133+
groupsErr: errors.New("connection refused"),
134+
}
135+
136+
_, err := IsNodeResourceTopologyEnabled(discoveryClient)
137+
Expect(err).To(HaveOccurred())
138+
})
139+
140+
It("should leave the gate untouched when discovery fails", func() {
141+
SetNodeResourceTopologyEnabledForTest(true)
142+
DeferCleanup(SetNodeResourceTopologyEnabledForTest, false)
143+
144+
discoveryClient := &erroringDiscovery{
145+
DiscoveryInterface: fake.NewClientset().Discovery(),
146+
groupsErr: errors.New("connection refused"),
147+
}
148+
149+
Expect(SetNodeResourceTopologyFeatureGate(discoveryClient)).NotTo(Succeed())
150+
Expect(NodeResourceTopologyEnabled()).To(BeTrue())
151+
})
49152
})
50153
})

pkg/scheduler/cache/cache.go

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -113,7 +113,7 @@ func registerSchedulerPodInformer(informerFactory informers.SharedInformerFactor
113113
}
114114

115115
// New returns a Cache implementation.
116-
func New(schedulerCacheParams *SchedulerCacheParams) Cache {
116+
func New(schedulerCacheParams *SchedulerCacheParams) (Cache, error) {
117117
return newSchedulerCache(schedulerCacheParams)
118118
}
119119

@@ -164,7 +164,7 @@ type SchedulerCache struct {
164164
K8sClusterPodAffinityInfo
165165
}
166166

167-
func newSchedulerCache(schedulerCacheParams *SchedulerCacheParams) *SchedulerCache {
167+
func newSchedulerCache(schedulerCacheParams *SchedulerCacheParams) (*SchedulerCache, error) {
168168
sc := &SchedulerCache{
169169
schedulingNodePoolParams: schedulerCacheParams.NodePoolParams,
170170
restrictNodeScheduling: schedulerCacheParams.RestrictNodeScheduling,
@@ -196,13 +196,16 @@ func newSchedulerCache(schedulerCacheParams *SchedulerCacheParams) *SchedulerCac
196196
sc.informerFactory = informers.NewSharedInformerFactory(sc.kubeClient, 0)
197197
registerSchedulerPodInformer(sc.informerFactory)
198198
if err := setSchedulerPodTransform(sc.informerFactory.Core().V1().Pods().Informer()); err != nil {
199-
log.InfraLogger.Errorf("Failed to set scheduler pod transform: %v", err)
200-
return nil
199+
return nil, fmt.Errorf("failed to set scheduler pod transform: %w", err)
201200
}
202201
sc.kubeAiSchedulerInformerFactory = kubeaischedulerinfo.NewSharedInformerFactory(sc.kubeAiSchedulerClient, 0)
203202

204-
featuregates.SetDRAFeatureGate(schedulerCacheParams.DiscoveryClient)
205-
featuregates.SetNodeResourceTopologyFeatureGate(schedulerCacheParams.DiscoveryClient)
203+
if err := featuregates.SetDRAFeatureGate(schedulerCacheParams.DiscoveryClient); err != nil {
204+
return nil, fmt.Errorf("failed to determine dynamic resource allocation availability: %w", err)
205+
}
206+
if err := featuregates.SetNodeResourceTopologyFeatureGate(schedulerCacheParams.DiscoveryClient); err != nil {
207+
return nil, fmt.Errorf("failed to determine node resource topology availability: %w", err)
208+
}
206209
if featuregates.NodeResourceTopologyEnabled() && schedulerCacheParams.NRTClient != nil {
207210
sc.nrtInformerFactory = nrtinformers.NewSharedInformerFactory(schedulerCacheParams.NRTClient, 0)
208211
}
@@ -223,12 +226,11 @@ func newSchedulerCache(schedulerCacheParams *SchedulerCacheParams) *SchedulerCac
223226
sc.restrictNodeScheduling, &sc.K8sClusterPodAffinityInfo, sc.scheduleCSIStorage, sc.fullHierarchyFairness, sc.StatusUpdater, sc.stuckInReleasingThreshold)
224227

225228
if err != nil {
226-
log.InfraLogger.Errorf("Failed to create cluster info object: %v", err)
227-
return nil
229+
return nil, fmt.Errorf("failed to create cluster info object: %w", err)
228230
}
229231
sc.clusterInfo = clusterInfo
230232

231-
return sc
233+
return sc, nil
232234
}
233235

234236
func (sc *SchedulerCache) Snapshot() (*api.ClusterInfo, error) {

pkg/scheduler/cache/cache_test.go

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -60,12 +60,13 @@ var _ = Describe("Cache", func() {
6060
Context("Pod informer filtering", func() {
6161
It("should filter failed pods while keeping succeeded pods, without filtering by scheduler name", func() {
6262
kubeClient := fake.NewSimpleClientset()
63-
cache := New(&SchedulerCacheParams{
63+
cache, err := New(&SchedulerCacheParams{
6464
KubeClient: kubeClient,
6565
KAISchedulerClient: kubeaischedulerfake.NewSimpleClientset(),
6666
NodePoolParams: &conf.SchedulingNodePoolParams{},
6767
DiscoveryClient: kubeClient.Discovery(),
6868
})
69+
Expect(err).NotTo(HaveOccurred())
6970

7071
stopCh := make(chan struct{})
7172
defer close(stopCh)
@@ -123,8 +124,9 @@ var _ = Describe("Cache", func() {
123124
DiscoveryClient: fakeClient.Discovery(),
124125
}
125126

126-
cache := New(params)
127+
cache, err := New(params)
127128

129+
Expect(err).NotTo(HaveOccurred())
128130
Expect(cache).NotTo(BeNil())
129131
Expect(featuregates.DynamicResourcesEnabled()).To(Equal(expectDRAAvailable))
130132
},
@@ -515,13 +517,14 @@ func setupCacheWithObjects(snapshot bool, objects []runtime.Object, kaiScheduler
515517
kubeClient := fake.NewSimpleClientset(objects...)
516518
kubeAiSchedulerClient := kubeaischedulerfake.NewSimpleClientset(kaiSchedulerObjects...)
517519

518-
cache := New(&SchedulerCacheParams{
520+
cache, err := New(&SchedulerCacheParams{
519521
KubeClient: kubeClient,
520522
KAISchedulerClient: kubeAiSchedulerClient,
521523
NodePoolParams: &conf.SchedulingNodePoolParams{},
522524
FullHierarchyFairness: true,
523525
DiscoveryClient: kubeClient.Discovery(),
524526
})
527+
Expect(err).NotTo(HaveOccurred())
525528

526529
stopCh := make(chan struct{})
527530
cache.Run(stopCh)

0 commit comments

Comments
 (0)