Skip to content

Commit 5afe170

Browse files
committed
Migrate from deprecated v1 Endpoints to discovery.k8s.io/v1 EndpointSlice
The Kubernetes v1 Endpoints API is deprecated in v1.33+. This migrates the Broker, InMemoryChannel, and EventTransform reconcilers to use EndpointSlice with label-based listing instead of name-based Endpoints lookups, and updates RBAC roles, test helpers, and all test files.
1 parent b3edbad commit 5afe170

28 files changed

Lines changed: 941 additions & 553 deletions

File tree

config/channels/in-memory-channel/roles/controller-clusterrole.yaml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -66,9 +66,9 @@ rules:
6666
- update
6767
- patch
6868
- apiGroups:
69-
- ""
69+
- "discovery.k8s.io"
7070
resources:
71-
- endpoints
71+
- endpointslices
7272
verbs:
7373
- get
7474
- list

config/core/roles/controller-clusterroles.yaml

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@ rules:
2727
- "secrets"
2828
- "configmaps"
2929
- "services"
30-
- "endpoints"
3130
- "events"
3231
- "serviceaccounts"
3332
- "pods"
@@ -41,6 +40,15 @@ rules:
4140
- "patch"
4241
- "watch"
4342

43+
- apiGroups:
44+
- "discovery.k8s.io"
45+
resources:
46+
- "endpointslices"
47+
verbs:
48+
- "get"
49+
- "list"
50+
- "watch"
51+
4452
# Brokers and the namespace annotation controllers manipulate Deployments.
4553
# RequestReply controller needs to manipulate StatefulSets
4654
- apiGroups:

pkg/apis/duck/lifecycle_helper.go

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ package duck
1818

1919
import (
2020
appsv1 "k8s.io/api/apps/v1"
21-
corev1 "k8s.io/api/core/v1"
21+
discoveryv1 "k8s.io/api/discovery/v1"
2222
)
2323

2424
// DeploymentIsAvailable determines if the provided deployment is available. Note that if it cannot
@@ -33,11 +33,14 @@ func DeploymentIsAvailable(d *appsv1.DeploymentStatus, def bool) bool {
3333
return def
3434
}
3535

36-
// EndpointsAreAvailable determines if the provided Endpoints are available.
37-
func EndpointsAreAvailable(ep *corev1.Endpoints) bool {
38-
for _, subset := range ep.Subsets {
39-
if len(subset.Addresses) > 0 {
40-
return true
36+
// EndpointSlicesAreAvailable determines if the provided EndpointSlices have any ready endpoints.
37+
func EndpointSlicesAreAvailable(epSlices []*discoveryv1.EndpointSlice) bool {
38+
for _, eps := range epSlices {
39+
for _, ep := range eps.Endpoints {
40+
// Per the K8s API spec, a nil Ready value should be interpreted as ready.
41+
if ep.Conditions.Ready == nil || *ep.Conditions.Ready {
42+
return true
43+
}
4144
}
4245
}
4346
return false

pkg/apis/eventing/v1/broker_lifecycle_mt.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ limitations under the License.
1717
package v1
1818

1919
import (
20-
corev1 "k8s.io/api/core/v1"
20+
discoveryv1 "k8s.io/api/discovery/v1"
2121

2222
"knative.dev/eventing/pkg/apis/duck"
2323
duckv1 "knative.dev/eventing/pkg/apis/duck/v1"
@@ -27,11 +27,11 @@ func (bs *BrokerStatus) MarkIngressFailed(reason, format string, args ...interfa
2727
bs.GetConditionSet().Manage(bs).MarkFalse(BrokerConditionIngress, reason, format, args...)
2828
}
2929

30-
func (bs *BrokerStatus) PropagateIngressAvailability(ep *corev1.Endpoints) {
31-
if duck.EndpointsAreAvailable(ep) {
30+
func (bs *BrokerStatus) PropagateIngressAvailability(epSlices []*discoveryv1.EndpointSlice) {
31+
if duck.EndpointSlicesAreAvailable(epSlices) {
3232
bs.GetConditionSet().Manage(bs).MarkTrue(BrokerConditionIngress)
3333
} else {
34-
bs.MarkIngressFailed("EndpointsUnavailable", "Endpoints %q are unavailable.", ep.Name)
34+
bs.MarkIngressFailed("EndpointSlicesUnavailable", "EndpointSlices are unavailable.")
3535
}
3636
}
3737

@@ -57,10 +57,10 @@ func (bs *BrokerStatus) MarkFilterFailed(reason, format string, args ...interfac
5757
bs.GetConditionSet().Manage(bs).MarkFalse(BrokerConditionFilter, reason, format, args...)
5858
}
5959

60-
func (bs *BrokerStatus) PropagateFilterAvailability(ep *corev1.Endpoints) {
61-
if duck.EndpointsAreAvailable(ep) {
60+
func (bs *BrokerStatus) PropagateFilterAvailability(epSlices []*discoveryv1.EndpointSlice) {
61+
if duck.EndpointSlicesAreAvailable(epSlices) {
6262
bs.GetConditionSet().Manage(bs).MarkTrue(BrokerConditionFilter)
6363
} else {
64-
bs.MarkFilterFailed("EndpointsUnavailable", "Endpoints %q are unavailable.", ep.Name)
64+
bs.MarkFilterFailed("EndpointSlicesUnavailable", "EndpointSlices are unavailable.")
6565
}
6666
}

pkg/apis/eventing/v1/broker_lifecycle_test.go

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import (
1818

1919
"github.com/google/go-cmp/cmp"
2020
corev1 "k8s.io/api/core/v1"
21+
discoveryv1 "k8s.io/api/discovery/v1"
2122
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2223
"k8s.io/utils/pointer"
2324

@@ -499,13 +500,13 @@ func TestBrokerIsReady(t *testing.T) {
499500
t.Run(test.name, func(t *testing.T) {
500501
bs := BrokerStatus{}
501502
if test.markIngressReady != nil {
502-
var ep *corev1.Endpoints
503+
var epSlices []*discoveryv1.EndpointSlice
503504
if *test.markIngressReady {
504-
ep = TestHelper.AvailableEndpoints()
505+
epSlices = TestHelper.AvailableEndpointSlices()
505506
} else {
506-
ep = TestHelper.UnavailableEndpoints()
507+
epSlices = TestHelper.UnavailableEndpointSlices()
507508
}
508-
bs.PropagateIngressAvailability(ep)
509+
bs.PropagateIngressAvailability(epSlices)
509510
}
510511
if test.markTriggerChannelReady != nil {
511512
var c *eventingduckv1.ChannelableStatus
@@ -534,13 +535,13 @@ func TestBrokerIsReady(t *testing.T) {
534535
}
535536

536537
if test.markFilterReady != nil {
537-
var ep *corev1.Endpoints
538+
var epSlices []*discoveryv1.EndpointSlice
538539
if *test.markFilterReady {
539-
ep = TestHelper.AvailableEndpoints()
540+
epSlices = TestHelper.AvailableEndpointSlices()
540541
} else {
541-
ep = TestHelper.UnavailableEndpoints()
542+
epSlices = TestHelper.UnavailableEndpointSlices()
542543
}
543-
bs.PropagateFilterAvailability(ep)
544+
bs.PropagateFilterAvailability(epSlices)
544545
}
545546

546547
if test.markAddressable == nil && test.address == nil {

pkg/apis/eventing/v1/test_helper.go

Lines changed: 17 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ package v1
1818

1919
import (
2020
corev1 "k8s.io/api/core/v1"
21+
discoveryv1 "k8s.io/api/discovery/v1"
2122

2223
"knative.dev/pkg/apis"
2324
duckv1 "knative.dev/pkg/apis/duck/v1"
@@ -59,9 +60,9 @@ func (testHelper) ReadySubscriptionStatus() *messagingv1.SubscriptionStatus {
5960

6061
func (t testHelper) ReadyBrokerStatus() *BrokerStatus {
6162
bs := &BrokerStatus{}
62-
bs.PropagateIngressAvailability(t.AvailableEndpoints())
63+
bs.PropagateIngressAvailability(t.AvailableEndpointSlices())
6364
bs.PropagateTriggerChannelReadiness(t.ReadyChannelStatus())
64-
bs.PropagateFilterAvailability(t.AvailableEndpoints())
65+
bs.PropagateFilterAvailability(t.AvailableEndpointSlices())
6566
bs.SetAddress(&duckv1.Addressable{
6667
URL: apis.HTTP("example.com"),
6768
})
@@ -72,9 +73,9 @@ func (t testHelper) ReadyBrokerStatus() *BrokerStatus {
7273

7374
func (t testHelper) ReadyBrokerStatusWithoutDLS() *BrokerStatus {
7475
bs := &BrokerStatus{}
75-
bs.PropagateIngressAvailability(t.AvailableEndpoints())
76+
bs.PropagateIngressAvailability(t.AvailableEndpointSlices())
7677
bs.PropagateTriggerChannelReadiness(t.ReadyChannelStatus())
77-
bs.PropagateFilterAvailability(t.AvailableEndpoints())
78+
bs.PropagateFilterAvailability(t.AvailableEndpointSlices())
7879
bs.SetAddress(&duckv1.Addressable{
7980
URL: apis.HTTP("example.com"),
8081
})
@@ -102,26 +103,24 @@ func (testHelper) FalseBrokerStatus() *BrokerStatus {
102103
return bs
103104
}
104105

105-
func (testHelper) UnavailableEndpoints() *corev1.Endpoints {
106-
ep := &corev1.Endpoints{}
107-
ep.Name = "unavailable"
108-
ep.Subsets = []corev1.EndpointSubset{{
109-
NotReadyAddresses: []corev1.EndpointAddress{{
110-
IP: "127.0.0.1",
106+
func (testHelper) UnavailableEndpointSlices() []*discoveryv1.EndpointSlice {
107+
ready := false
108+
return []*discoveryv1.EndpointSlice{{
109+
Endpoints: []discoveryv1.Endpoint{{
110+
Addresses: []string{"127.0.0.1"},
111+
Conditions: discoveryv1.EndpointConditions{Ready: &ready},
111112
}},
112113
}}
113-
return ep
114114
}
115115

116-
func (testHelper) AvailableEndpoints() *corev1.Endpoints {
117-
ep := &corev1.Endpoints{}
118-
ep.Name = "available"
119-
ep.Subsets = []corev1.EndpointSubset{{
120-
Addresses: []corev1.EndpointAddress{{
121-
IP: "127.0.0.1",
116+
func (testHelper) AvailableEndpointSlices() []*discoveryv1.EndpointSlice {
117+
ready := true
118+
return []*discoveryv1.EndpointSlice{{
119+
Endpoints: []discoveryv1.Endpoint{{
120+
Addresses: []string{"127.0.0.1"},
121+
Conditions: discoveryv1.EndpointConditions{Ready: &ready},
122122
}},
123123
}}
124-
return ep
125124
}
126125

127126
func (testHelper) ReadyChannelStatus() *eventingduckv1.ChannelableStatus {

pkg/reconciler/broker/broker.go

Lines changed: 16 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424

2525
"go.uber.org/zap"
2626
corev1 "k8s.io/api/core/v1"
27+
discoveryv1 "k8s.io/api/discovery/v1"
2728
"k8s.io/apimachinery/pkg/api/equality"
2829
apierrs "k8s.io/apimachinery/pkg/api/errors"
2930
"k8s.io/apimachinery/pkg/api/meta"
@@ -33,6 +34,7 @@ import (
3334
"k8s.io/apimachinery/pkg/types"
3435
"k8s.io/client-go/dynamic"
3536
corev1listers "k8s.io/client-go/listers/core/v1"
37+
discoveryv1listers "k8s.io/client-go/listers/discovery/v1"
3638
"k8s.io/utils/pointer"
3739
"knative.dev/pkg/kmeta"
3840
"knative.dev/pkg/resolver"
@@ -72,10 +74,10 @@ type Reconciler struct {
7274
dynamicClientSet dynamic.Interface
7375

7476
// listers index properties about resources
75-
endpointsLister corev1listers.EndpointsLister
76-
subscriptionLister messaginglisters.SubscriptionLister
77-
configmapLister corev1listers.ConfigMapLister
78-
secretLister corev1listers.SecretLister
77+
endpointSliceLister discoveryv1listers.EndpointSliceLister
78+
subscriptionLister messaginglisters.SubscriptionLister
79+
configmapLister corev1listers.ConfigMapLister
80+
secretLister corev1listers.SecretLister
7981

8082
channelableTracker ducklib.ListableTracker
8183

@@ -188,21 +190,25 @@ func (r *Reconciler) ReconcileKind(ctx context.Context, b *eventingv1.Broker) pk
188190

189191
b.Status.PropagateTriggerChannelReadiness(channelStatus)
190192

191-
filterEndpoints, err := r.endpointsLister.Endpoints(system.Namespace()).Get(names.BrokerFilterName)
193+
filterEpSlices, err := r.endpointSliceLister.EndpointSlices(system.Namespace()).List(labels.SelectorFromSet(labels.Set{
194+
discoveryv1.LabelServiceName: names.BrokerFilterName,
195+
}))
192196
if err != nil {
193-
logging.FromContext(ctx).Errorw("Problem getting endpoints for filter", zap.String("namespace", system.Namespace()), zap.Error(err))
197+
logging.FromContext(ctx).Errorw("Problem getting EndpointSlices for filter", zap.String("namespace", system.Namespace()), zap.Error(err))
194198
b.Status.MarkFilterFailed("ServiceFailure", "%v", err)
195199
return err
196200
}
197-
b.Status.PropagateFilterAvailability(filterEndpoints)
201+
b.Status.PropagateFilterAvailability(filterEpSlices)
198202

199-
ingressEndpoints, err := r.endpointsLister.Endpoints(system.Namespace()).Get(names.BrokerIngressName)
203+
ingressEpSlices, err := r.endpointSliceLister.EndpointSlices(system.Namespace()).List(labels.SelectorFromSet(labels.Set{
204+
discoveryv1.LabelServiceName: names.BrokerIngressName,
205+
}))
200206
if err != nil {
201-
logging.FromContext(ctx).Errorw("Problem getting endpoints for ingress", zap.String("namespace", system.Namespace()), zap.Error(err))
207+
logging.FromContext(ctx).Errorw("Problem getting EndpointSlices for ingress", zap.String("namespace", system.Namespace()), zap.Error(err))
202208
b.Status.MarkIngressFailed("ServiceFailure", "%v", err)
203209
return err
204210
}
205-
b.Status.PropagateIngressAvailability(ingressEndpoints)
211+
b.Status.PropagateIngressAvailability(ingressEpSlices)
206212

207213
if b.Spec.Delivery != nil && b.Spec.Delivery.DeadLetterSink != nil {
208214
deadLetterSinkAddr, err := r.uriResolver.AddressableFromDestinationV1(ctx, *b.Spec.Delivery.DeadLetterSink, b)

0 commit comments

Comments
 (0)