Skip to content

Commit 2a59266

Browse files
committed
Fixed admission and scheduler PodDisruptionBudget reconciliation
Signed-off-by: dttung2905 <ttdao.2015@accountancy.smu.edu.sg>
1 parent 64f3e37 commit 2a59266

7 files changed

Lines changed: 486 additions & 0 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: Fixed admission and scheduler PodDisruptionBudget reconciliation, including drift recovery and cleanup when no longer desired.
3+
time: 2026-07-18T21:29:43.130860557+01:00
4+
custom:
5+
Author: dttung2905
6+
Issue: "1477"

pkg/operator/controller/integration_tests/config_controller_test.go

Lines changed: 169 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,14 +8,18 @@ import (
88
"os"
99

1010
kaiv1 "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1"
11+
kaiv1admission "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1/admission"
1112
kaiv1binder "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1/binder"
13+
kaiv1scheduler "github.com/kai-scheduler/KAI-scheduler/pkg/apis/kai/v1/scheduler"
1214
"github.com/kai-scheduler/KAI-scheduler/pkg/common/constants"
1315
. "github.com/onsi/ginkgo/v2"
1416
. "github.com/onsi/gomega"
1517

1618
nvidiav1 "github.com/kai-scheduler/KAI-scheduler/third_party/nvidia/gpu-operator/api/nvidia/v1"
1719

1820
appsv1 "k8s.io/api/apps/v1"
21+
policyv1 "k8s.io/api/policy/v1"
22+
apierrors "k8s.io/apimachinery/pkg/api/errors"
1923
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2024
"k8s.io/apimachinery/pkg/types"
2125
"k8s.io/utils/ptr"
@@ -116,4 +120,169 @@ var _ = Describe("KAIConfigController", Ordered, func() {
116120
})
117121
})
118122

123+
Context("Reconciling the admission PodDisruptionBudget", Ordered, func() {
124+
const pdbName = "admission"
125+
126+
updateAdmissionPDBConfig := func(ctx context.Context, replicas int32, enabled bool) {
127+
Eventually(func() error {
128+
currentConfig := &kaiv1.Config{}
129+
if err := k8sClient.Get(
130+
ctx,
131+
types.NamespacedName{Name: constants.DefaultKAIConfigSingeltonInstanceName},
132+
currentConfig,
133+
); err != nil {
134+
return err
135+
}
136+
currentConfig.Spec.Admission = &kaiv1admission.Admission{
137+
Replicas: ptr.To(replicas),
138+
Service: &v1common.Service{
139+
Enabled: ptr.To(true),
140+
PodDisruptionBudget: &v1common.PodDisruptionBudget{
141+
Enabled: ptr.To(enabled),
142+
MaxUnavailable: ptr.To(int32(1)),
143+
},
144+
},
145+
}
146+
return k8sClient.Update(ctx, currentConfig)
147+
}, "10s", "200ms").Should(Succeed())
148+
}
149+
150+
getPDB := func(ctx context.Context) (*policyv1.PodDisruptionBudget, error) {
151+
pdb := &policyv1.PodDisruptionBudget{}
152+
err := k8sClient.Get(ctx, types.NamespacedName{
153+
Name: pdbName,
154+
Namespace: kaiConfig.Spec.Namespace,
155+
}, pdb)
156+
return pdb, err
157+
}
158+
159+
It("creates and watches the PDB", func(ctx context.Context) {
160+
updateAdmissionPDBConfig(ctx, 2, true)
161+
162+
Eventually(func(g Gomega) {
163+
pdb, err := getPDB(ctx)
164+
g.Expect(err).NotTo(HaveOccurred())
165+
g.Expect(metav1.GetControllerOf(pdb)).NotTo(BeNil())
166+
g.Expect(metav1.GetControllerOf(pdb).Kind).To(Equal("Config"))
167+
}, "10s", "200ms").Should(Succeed())
168+
169+
pdb, err := getPDB(ctx)
170+
Expect(err).NotTo(HaveOccurred())
171+
Expect(k8sClient.Delete(ctx, pdb)).To(Succeed())
172+
173+
Eventually(func() error {
174+
_, err := getPDB(ctx)
175+
return err
176+
}, "10s", "200ms").Should(Succeed())
177+
})
178+
179+
It("removes the PDB after scaling admission down to one replica", func(ctx context.Context) {
180+
updateAdmissionPDBConfig(ctx, 1, true)
181+
182+
Eventually(func() bool {
183+
_, err := getPDB(ctx)
184+
return apierrors.IsNotFound(err)
185+
}, "10s", "200ms").Should(BeTrue())
186+
})
187+
188+
It("removes the PDB when it is disabled", func(ctx context.Context) {
189+
updateAdmissionPDBConfig(ctx, 2, true)
190+
Eventually(func() error {
191+
_, err := getPDB(ctx)
192+
return err
193+
}, "10s", "200ms").Should(Succeed())
194+
195+
updateAdmissionPDBConfig(ctx, 2, false)
196+
Eventually(func() bool {
197+
_, err := getPDB(ctx)
198+
return apierrors.IsNotFound(err)
199+
}, "10s", "200ms").Should(BeTrue())
200+
})
201+
})
202+
203+
Context("Reconciling the scheduler PodDisruptionBudget", Ordered, func() {
204+
const (
205+
shardName = "pdb-test"
206+
pdbName = "kai-scheduler-" + shardName
207+
)
208+
209+
updateSchedulerPDBConfig := func(ctx context.Context, replicas int32, enabled bool) {
210+
Eventually(func() error {
211+
currentConfig := &kaiv1.Config{}
212+
if err := k8sClient.Get(
213+
ctx,
214+
types.NamespacedName{Name: constants.DefaultKAIConfigSingeltonInstanceName},
215+
currentConfig,
216+
); err != nil {
217+
return err
218+
}
219+
currentConfig.Spec.Scheduler = &kaiv1scheduler.Scheduler{
220+
Replicas: ptr.To(replicas),
221+
Service: &v1common.Service{
222+
Enabled: ptr.To(true),
223+
PodDisruptionBudget: &v1common.PodDisruptionBudget{
224+
Enabled: ptr.To(enabled),
225+
MaxUnavailable: ptr.To(int32(1)),
226+
},
227+
},
228+
}
229+
return k8sClient.Update(ctx, currentConfig)
230+
}, "10s", "200ms").Should(Succeed())
231+
}
232+
233+
getPDB := func(ctx context.Context) (*policyv1.PodDisruptionBudget, error) {
234+
pdb := &policyv1.PodDisruptionBudget{}
235+
err := k8sClient.Get(ctx, types.NamespacedName{
236+
Name: pdbName,
237+
Namespace: kaiConfig.Spec.Namespace,
238+
}, pdb)
239+
return pdb, err
240+
}
241+
242+
BeforeAll(func(ctx context.Context) {
243+
updateSchedulerPDBConfig(ctx, 2, true)
244+
Expect(k8sClient.Create(ctx, &kaiv1.SchedulingShard{
245+
ObjectMeta: metav1.ObjectMeta{Name: shardName},
246+
})).To(Succeed())
247+
})
248+
249+
AfterAll(func(ctx context.Context) {
250+
shard := &kaiv1.SchedulingShard{}
251+
err := k8sClient.Get(ctx, types.NamespacedName{Name: shardName}, shard)
252+
if err == nil {
253+
Expect(k8sClient.Delete(ctx, shard)).To(Succeed())
254+
} else {
255+
Expect(apierrors.IsNotFound(err)).To(BeTrue())
256+
}
257+
})
258+
259+
It("creates and watches one PDB for the shard", func(ctx context.Context) {
260+
Eventually(func(g Gomega) {
261+
pdb, err := getPDB(ctx)
262+
g.Expect(err).NotTo(HaveOccurred())
263+
g.Expect(metav1.GetControllerOf(pdb)).NotTo(BeNil())
264+
g.Expect(metav1.GetControllerOf(pdb).Kind).To(Equal("SchedulingShard"))
265+
g.Expect(metav1.GetControllerOf(pdb).Name).To(Equal(shardName))
266+
}, "10s", "200ms").Should(Succeed())
267+
268+
pdb, err := getPDB(ctx)
269+
Expect(err).NotTo(HaveOccurred())
270+
Expect(k8sClient.Delete(ctx, pdb)).To(Succeed())
271+
272+
Eventually(func() error {
273+
_, err := getPDB(ctx)
274+
return err
275+
}, "10s", "200ms").Should(Succeed())
276+
})
277+
278+
It("removes the shard PDB after scheduler scales down", func(ctx context.Context) {
279+
updateSchedulerPDBConfig(ctx, 1, true)
280+
281+
Eventually(func() bool {
282+
_, err := getPDB(ctx)
283+
return apierrors.IsNotFound(err)
284+
}, "10s", "200ms").Should(BeTrue())
285+
})
286+
})
287+
119288
})

pkg/operator/controller/integration_tests/suite_test.go

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,10 @@ import (
3333
monitoringv1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1"
3434
admissionv1 "k8s.io/api/admissionregistration/v1"
3535
appsv1 "k8s.io/api/apps/v1"
36+
coordinationv1 "k8s.io/api/coordination/v1"
3637
v1 "k8s.io/api/core/v1"
38+
discoveryv1 "k8s.io/api/discovery/v1"
39+
policyv1 "k8s.io/api/policy/v1"
3740
apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
3841
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
3942
"k8s.io/apimachinery/pkg/runtime"
@@ -73,6 +76,9 @@ var _ = BeforeSuite(func() {
7376
Expect(appsv1.AddToScheme(scheme)).To(Succeed())
7477
Expect(apiserver.AddToScheme(scheme)).To(Succeed())
7578
Expect(admissionv1.AddToScheme(scheme)).To(Succeed())
79+
Expect(coordinationv1.AddToScheme(scheme)).To(Succeed())
80+
Expect(discoveryv1.AddToScheme(scheme)).To(Succeed())
81+
Expect(policyv1.AddToScheme(scheme)).To(Succeed())
7682
Expect(apiextensionsv1.AddToScheme(scheme)).To(Succeed())
7783

7884
By("bootstrapping test environment")
@@ -125,6 +131,14 @@ var _ = BeforeSuite(func() {
125131
err = configReconciler.SetupWithManager(k8sManager)
126132
Expect(err).NotTo(HaveOccurred())
127133

134+
shardReconciler := controller.NewSchedulingShardReconciler(
135+
k8sManager.GetClient(),
136+
k8sManager.GetScheme(),
137+
)
138+
shardReconciler.SetOperands(controller.OperandsForShard)
139+
err = shardReconciler.SetupWithManager(k8sManager)
140+
Expect(err).NotTo(HaveOccurred())
141+
128142
go func() {
129143
defer GinkgoRecover()
130144
err = k8sManager.Start(testContext)

pkg/operator/operands/deployable/deployable_test.go

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@ import (
1919
monitoringv1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1"
2020
appsv1 "k8s.io/api/apps/v1"
2121
v1 "k8s.io/api/core/v1"
22+
policyv1 "k8s.io/api/policy/v1"
2223
apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
24+
apierrors "k8s.io/apimachinery/pkg/api/errors"
2325
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2426
vpav1 "k8s.io/autoscaler/vertical-pod-autoscaler/pkg/apis/autoscaling.k8s.io/v1"
2527
"k8s.io/client-go/kubernetes/scheme"
@@ -54,6 +56,7 @@ var _ = Describe("Deployable", func() {
5456
Expect(kaiv1.AddToScheme(testScheme)).To(Succeed())
5557
Expect(apiextensionsv1.AddToScheme(testScheme)).To(Succeed())
5658
Expect(monitoringv1.AddToScheme(testScheme)).To(Succeed())
59+
Expect(policyv1.AddToScheme(testScheme)).To(Succeed())
5760
Expect(vpav1.AddToScheme(testScheme)).To(Succeed())
5861

5962
fakeClientBuilder = fake.NewClientBuilder().
@@ -215,6 +218,77 @@ var _ = Describe("Deployable", func() {
215218
Expect(updateCalls).To(Equal(1))
216219
})
217220
})
221+
222+
Context("PodDisruptionBudget lifecycle", func() {
223+
var existingPDB *policyv1.PodDisruptionBudget
224+
225+
BeforeEach(func() {
226+
existingPDB = &policyv1.PodDisruptionBudget{
227+
TypeMeta: metav1.TypeMeta{
228+
Kind: "PodDisruptionBudget",
229+
APIVersion: policyv1.SchemeGroupVersion.String(),
230+
},
231+
ObjectMeta: metav1.ObjectMeta{
232+
Name: "admission",
233+
Namespace: "kai-scheduler",
234+
OwnerReferences: []metav1.OwnerReference{{
235+
APIVersion: kaiv1.GroupVersion.String(),
236+
Kind: "Config",
237+
Name: kaiConfig.Name,
238+
UID: kaiConfig.UID,
239+
Controller: ptr.To(true),
240+
}},
241+
},
242+
}
243+
})
244+
245+
It("collects an existing PDB instead of attempting to create it again", func() {
246+
pdbCreateCalls := 0
247+
builder := fakeClientBuilder.
248+
WithObjects(existingPDB).
249+
WithInterceptorFuncs(interceptor.Funcs{
250+
Create: func(
251+
ctx context.Context,
252+
runtimeClient client.WithWatch,
253+
obj client.Object,
254+
opts ...client.CreateOption,
255+
) error {
256+
if _, ok := obj.(*policyv1.PodDisruptionBudget); ok {
257+
pdbCreateCalls++
258+
}
259+
return runtimeClient.Create(ctx, obj, opts...)
260+
},
261+
})
262+
fakeClient := getFakeClient(builder, known_types.KAIConfigRegisteredCollectible)
263+
deployable := New(
264+
[]operands.Operand{&fakePDBOperand{enabled: true}},
265+
known_types.KAIConfigRegisteredCollectible,
266+
)
267+
268+
Expect(deployable.Deploy(context.Background(), fakeClient, kaiConfig, kaiConfig)).To(Succeed())
269+
Expect(pdbCreateCalls).To(BeZero())
270+
})
271+
272+
It("deletes an owned PDB when it is no longer desired", func() {
273+
fakeClient := getFakeClient(
274+
fakeClientBuilder.WithObjects(existingPDB),
275+
known_types.KAIConfigRegisteredCollectible,
276+
)
277+
deployable := New(
278+
[]operands.Operand{&fakePDBOperand{enabled: false}},
279+
known_types.KAIConfigRegisteredCollectible,
280+
)
281+
282+
Expect(deployable.Deploy(context.Background(), fakeClient, kaiConfig, kaiConfig)).To(Succeed())
283+
284+
err := fakeClient.Get(
285+
context.Background(),
286+
client.ObjectKeyFromObject(existingPDB),
287+
&policyv1.PodDisruptionBudget{},
288+
)
289+
Expect(apierrors.IsNotFound(err)).To(BeTrue())
290+
})
291+
})
218292
})
219293

220294
Describe("IsDeployed", func() {
@@ -442,6 +516,53 @@ type fakeOperand struct {
442516
name string
443517
}
444518

519+
type fakePDBOperand struct {
520+
enabled bool
521+
}
522+
523+
func (f *fakePDBOperand) DesiredState(
524+
ctx context.Context,
525+
runtimeClient client.Reader,
526+
_ *kaiv1.Config,
527+
) ([]client.Object, error) {
528+
if !f.enabled {
529+
return nil, nil
530+
}
531+
532+
pdb := &policyv1.PodDisruptionBudget{}
533+
err := runtimeClient.Get(ctx, client.ObjectKey{Name: "admission", Namespace: "kai-scheduler"}, pdb)
534+
if err != nil && !apierrors.IsNotFound(err) {
535+
return nil, err
536+
}
537+
pdb.TypeMeta = metav1.TypeMeta{
538+
Kind: "PodDisruptionBudget",
539+
APIVersion: policyv1.SchemeGroupVersion.String(),
540+
}
541+
pdb.Name = "admission"
542+
pdb.Namespace = "kai-scheduler"
543+
return []client.Object{pdb}, nil
544+
}
545+
546+
func (f *fakePDBOperand) IsDeployed(context.Context, client.Reader) (bool, error) {
547+
return true, nil
548+
}
549+
550+
func (f *fakePDBOperand) IsAvailable(context.Context, client.Reader) (bool, error) {
551+
return true, nil
552+
}
553+
554+
func (f *fakePDBOperand) Monitor(context.Context, client.Reader, *kaiv1.Config) error {
555+
return nil
556+
}
557+
558+
func (f *fakePDBOperand) HasMissingDependencies(context.Context, client.Reader, *kaiv1.Config) (string, error) {
559+
return "", nil
560+
}
561+
562+
func (f *fakePDBOperand) Name() string {
563+
return "fakePDBOperand"
564+
}
565+
445566
func (f *fakeOperand) DesiredState(_ context.Context, _ client.Reader, _ *kaiv1.Config) ([]client.Object, error) {
446567
return []client.Object{
447568
&v1.Pod{

pkg/operator/operands/known_types/known_types.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ func init() {
4949
registerCustomResourceDefinitions()
5050
registerPrometheus()
5151
registerVerticalPodAutoscalers()
52+
registerPodDisruptionBudgets()
5253
}
5354

5455
func SetupKAIConfigOwned(fn *Collectable) {

0 commit comments

Comments
 (0)