Skip to content

Commit 751dc5e

Browse files
committed
- Remove zone filtering from dataclient - return all endpoints with zone metadata
- Ensure route construction always creates LBEndpoints with zone labels - Implement threshold-based zone filtering in routesrv when serving routes - Update expectations of testing zone-aware-routing in dataclient Signed-off-by: greeshma1196 <greeshma.mathew@gmail.com>
1 parent e0f7bc9 commit 751dc5e

11 files changed

Lines changed: 163 additions & 115 deletions

File tree

dataclients/kubernetes/clusterclient.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -677,6 +677,7 @@ func (c *clusterClient) fetchClusterState() (*clusterState, error) {
677677
routeGroups: routeGroups,
678678
services: services,
679679
cachedEndpoints: make(map[endpointID][]string),
680+
cachedEndpointSlices: make(map[endpointID][]skipperEndpoint),
680681
enableEndpointSlices: c.enableEndpointSlices,
681682
}
682683

dataclients/kubernetes/clusterstate.go

Lines changed: 71 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ type clusterState struct {
1919
endpointSlices map[definitions.ResourceID]*skipperEndpointSlice
2020
secrets map[definitions.ResourceID]*secret
2121
cachedEndpoints map[endpointID][]string
22+
cachedEndpointSlices map[endpointID][]skipperEndpoint
2223
enableEndpointSlices bool
2324
}
2425

@@ -49,8 +50,8 @@ func (state *clusterState) getServiceRG(namespace, name string) (*service, error
4950
return s, nil
5051
}
5152

52-
// GetEndpointsByService returns the skipper endpoints for kubernetes endpoints or endpointslices.
53-
func (state *clusterState) GetEndpointsByService(zone, namespace, name, protocol string, servicePort *servicePort) ([]string, bool) {
53+
// GetEndpointsByService returns the skipper endpoints for kubernetes endpoints.
54+
func (state *clusterState) GetEndpointsByService(namespace, name, protocol string, servicePort *servicePort) []string {
5455
epID := endpointID{
5556
ResourceID: newResourceID(namespace, name),
5657
Protocol: protocol,
@@ -60,34 +61,48 @@ func (state *clusterState) GetEndpointsByService(zone, namespace, name, protocol
6061
state.mu.Lock()
6162
defer state.mu.Unlock()
6263
if cached, ok := state.cachedEndpoints[epID]; ok {
63-
return cached, false
64+
return cached
6465
}
6566

6667
var targets []string
67-
var targetsByZone []string
68-
if state.enableEndpointSlices {
69-
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
70-
targets, targetsByZone = eps.targetsByServicePort(zone, "TCP", protocol, servicePort)
71-
} else {
72-
return nil, false
73-
}
68+
if ep, ok := state.endpoints[epID.ResourceID]; ok {
69+
targets = ep.targetsByServicePort(protocol, servicePort)
7470
} else {
75-
if ep, ok := state.endpoints[epID.ResourceID]; ok {
76-
targets = ep.targetsByServicePort(protocol, servicePort)
77-
} else {
78-
return nil, false
79-
}
80-
}
81-
82-
if len(targetsByZone) >= minEndpointsByZone {
83-
sort.Strings(targetsByZone)
84-
state.cachedEndpoints[epID] = targetsByZone
85-
return targetsByZone, true
71+
return nil
8672
}
8773

8874
sort.Strings(targets)
8975
state.cachedEndpoints[epID] = targets
90-
return targets, false
76+
return targets
77+
}
78+
79+
// GetEndpointSlicesByService returns the skipper endpointslices for kubernetes endpointslices.
80+
func (state *clusterState) GetEndpointSlicesByService(namespace, name, protocol string, servicePort *servicePort) []skipperEndpoint {
81+
epID := endpointID{
82+
ResourceID: newResourceID(namespace, name),
83+
Protocol: protocol,
84+
TargetPort: servicePort.TargetPort.String(),
85+
}
86+
87+
state.mu.Lock()
88+
defer state.mu.Unlock()
89+
if cached, ok := state.cachedEndpointSlices[epID]; ok {
90+
return cached
91+
}
92+
93+
var targets []skipperEndpoint
94+
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
95+
targets = eps.targetsByServicePort("TCP", protocol, servicePort)
96+
} else {
97+
return nil
98+
}
99+
100+
sort.Slice(targets, func(i, j int) bool {
101+
return targets[i].Address < targets[j].Address
102+
})
103+
104+
state.cachedEndpointSlices[epID] = targets
105+
return targets
91106
}
92107

93108
// getEndpointAddresses returns the list of all addresses for the given service using endpoints or endpointslices.
@@ -121,7 +136,7 @@ func (state *clusterState) getEndpointAddresses(zone, namespace, name string) []
121136
}
122137

123138
// GetEndpointsByTarget returns the skipper endpoints for kubernetes endpoints or endpointslices.
124-
func (state *clusterState) GetEndpointsByTarget(zone, namespace, name, protocol, scheme string, target *definitions.BackendPort) ([]string, bool) {
139+
func (state *clusterState) GetEndpointsByTarget(namespace, name, protocol, scheme string, target *definitions.BackendPort) []string {
125140
epID := endpointID{
126141
ResourceID: newResourceID(namespace, name),
127142
Protocol: protocol,
@@ -131,32 +146,45 @@ func (state *clusterState) GetEndpointsByTarget(zone, namespace, name, protocol,
131146
state.mu.Lock()
132147
defer state.mu.Unlock()
133148
if cached, ok := state.cachedEndpoints[epID]; ok {
134-
return cached, false
149+
return cached
135150
}
136151

137152
var targets []string
138-
var targetsByZone []string
139-
if state.enableEndpointSlices {
140-
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
141-
targets, targetsByZone = eps.targetsByServiceTarget(zone, protocol, scheme, target)
142-
} else {
143-
return nil, false
144-
}
145-
} else {
146-
if ep, ok := state.endpoints[epID.ResourceID]; ok {
147-
targets = ep.targetsByServiceTarget(scheme, target)
148-
} else {
149-
return nil, false
150-
}
151-
}
152153

153-
if len(targetsByZone) >= minEndpointsByZone {
154-
sort.Strings(targetsByZone)
155-
state.cachedEndpoints[epID] = targetsByZone
156-
return targetsByZone, true
154+
if ep, ok := state.endpoints[epID.ResourceID]; ok {
155+
targets = ep.targetsByServiceTarget(scheme, target)
156+
} else {
157+
return nil
157158
}
158159

159160
sort.Strings(targets)
160161
state.cachedEndpoints[epID] = targets
161-
return targets, false
162+
return targets
163+
}
164+
165+
func (state *clusterState) GetEndpointSlicesByTarget(namespace, name, protocol, scheme string, target *definitions.BackendPort) []skipperEndpoint {
166+
epID := endpointID{
167+
ResourceID: newResourceID(namespace, name),
168+
Protocol: protocol,
169+
TargetPort: target.String(),
170+
}
171+
172+
state.mu.Lock()
173+
defer state.mu.Unlock()
174+
if cached, ok := state.cachedEndpointSlices[epID]; ok {
175+
return cached
176+
}
177+
178+
var targets []skipperEndpoint
179+
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
180+
targets = eps.targetsByServiceTarget(protocol, scheme, target)
181+
} else {
182+
return nil
183+
}
184+
185+
sort.Slice(targets, func(i, j int) bool {
186+
return targets[i].Address < targets[j].Address
187+
})
188+
state.cachedEndpointSlices[epID] = targets
189+
return targets
162190
}

dataclients/kubernetes/clusterstate_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ func benchmarkCachedEndpoints(b *testing.B, n int) {
5757
b.ResetTimer()
5858
dummy := []string{}
5959
for i := 0; i < b.N; i++ {
60-
dummy, _ = cs.GetEndpointsByTarget("", "default", "foo-0", "TCP", "http", &definitions.BackendPort{})
60+
dummy = cs.GetEndpointsByTarget("default", "foo-0", "TCP", "http", &definitions.BackendPort{})
6161
}
6262
dummy2 = dummy
6363
}

dataclients/kubernetes/endpointslices.go

Lines changed: 10 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ func (eps *skipperEndpointSlice) getPort(protocol, pName string, pValue int) int
4343

4444
return port
4545
}
46-
func (eps *skipperEndpointSlice) targetsByServicePort(zone, protocol, scheme string, servicePort *servicePort) ([]string, []string) {
46+
func (eps *skipperEndpointSlice) targetsByServicePort(protocol, scheme string, servicePort *servicePort) []skipperEndpoint {
4747
var port int
4848
if servicePort.Name != "" {
4949
port = eps.getPort(protocol, servicePort.Name, servicePort.Port)
@@ -57,31 +57,25 @@ func (eps *skipperEndpointSlice) targetsByServicePort(zone, protocol, scheme str
5757
port = eps.getPort(protocol, servicePort.Name, servicePort.Port)
5858
}
5959

60-
result := make([]string, 0, len(eps.Endpoints))
61-
resultByZone := make([]string, 0, len(eps.Endpoints))
60+
result := make([]skipperEndpoint, 0, len(eps.Endpoints))
6261
for _, ep := range eps.Endpoints {
63-
if ep.Zone == zone {
64-
resultByZone = append(resultByZone, formatEndpointString(ep.Address, scheme, port))
65-
}
66-
result = append(result, formatEndpointString(ep.Address, scheme, port))
62+
addr := formatEndpointString(ep.Address, scheme, port)
63+
result = append(result, skipperEndpoint{Address: addr, Zone: ep.Zone})
6764
}
68-
return result, resultByZone
65+
return result
6966
}
7067

71-
func (eps *skipperEndpointSlice) targetsByServiceTarget(zone, protocol, scheme string, serviceTarget *definitions.BackendPort) ([]string, []string) {
68+
func (eps *skipperEndpointSlice) targetsByServiceTarget(protocol, scheme string, serviceTarget *definitions.BackendPort) []skipperEndpoint {
7269
pName, _ := serviceTarget.Value.(string)
7370
pValue, _ := serviceTarget.Value.(int)
7471
port := eps.getPort(protocol, pName, pValue)
7572

76-
result := make([]string, 0, len(eps.Endpoints))
77-
resultByZone := make([]string, 0, len(eps.Endpoints))
73+
result := make([]skipperEndpoint, 0, len(eps.Endpoints))
7874
for _, ep := range eps.Endpoints {
79-
if ep.Zone == zone {
80-
resultByZone = append(resultByZone, formatEndpointString(ep.Address, scheme, port))
81-
}
82-
result = append(result, formatEndpointString(ep.Address, scheme, port))
75+
addr := formatEndpointString(ep.Address, scheme, port)
76+
result = append(result, skipperEndpoint{Address: addr, Zone: ep.Zone})
8377
}
84-
return result, resultByZone
78+
return result
8579
}
8680

8781
func (eps *skipperEndpointSlice) addressesByZone(zone string) []string {

dataclients/kubernetes/ingressv1.go

Lines changed: 39 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -77,16 +77,15 @@ func convertPathRuleV1(
7777
ns := metadata.Namespace
7878
name := metadata.Name
7979

80-
zoneTarget := false
81-
8280
if prule.Backend == nil {
8381
return nil, fmt.Errorf("invalid path rule, missing backend in: %s/%s/%s", ns, name, host)
8482
}
8583

8684
var (
87-
eps []string
88-
err error
89-
svc *service
85+
eps []string
86+
epSlices []skipperEndpoint
87+
err error
88+
svc *service
9089
)
9190

9291
var hostRegexp []string
@@ -123,7 +122,14 @@ func convertPathRuleV1(
123122
protocol = p
124123
}
125124

126-
eps, zoneTarget = state.GetEndpointsByService(ic.zone, ns, svcName, protocol, servicePort)
125+
if state.enableEndpointSlices {
126+
epSlices = state.GetEndpointSlicesByService(ns, svcName, protocol, servicePort)
127+
for _, ep := range epSlices {
128+
eps = append(eps, ep.Address)
129+
}
130+
} else {
131+
eps = state.GetEndpointsByService(ns, svcName, protocol, servicePort)
132+
}
127133
}
128134
if len(eps) == 0 {
129135
ic.logger.Tracef("Target endpoints not found, shuntroute for %s:%s", svcName, svcPort)
@@ -161,10 +167,10 @@ func convertPathRuleV1(
161167
HostRegexps: hostRegexp,
162168
}
163169

164-
if zoneTarget {
170+
if state.enableEndpointSlices {
165171
var lbeps []*eskip.LBEndpoint
166-
for _, ep := range eps {
167-
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep, Zone: ic.zone})
172+
for _, ep := range epSlices {
173+
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep.Address, Zone: ep.Zone})
168174
}
169175
r.LBEndpoints = lbeps
170176
} else {
@@ -361,15 +367,14 @@ func (ing *ingress) convertDefaultBackendV1(
361367
}
362368

363369
var (
364-
eps []string
365-
err error
366-
ns = i.Metadata.Namespace
367-
name = i.Metadata.Name
368-
svcName = i.Spec.DefaultBackend.Service.Name
369-
svcPort = i.Spec.DefaultBackend.Service.Port
370+
eps []string
371+
epSlices []skipperEndpoint
372+
err error
373+
ns = i.Metadata.Namespace
374+
name = i.Metadata.Name
375+
svcName = i.Spec.DefaultBackend.Service.Name
376+
svcPort = i.Spec.DefaultBackend.Service.Port
370377
)
371-
zoneTarget := false
372-
dataclientZone := ing.zone
373378

374379
svc, err := state.getService(ns, svcName)
375380
if err != nil {
@@ -397,14 +402,20 @@ func (ing *ingress) convertDefaultBackendV1(
397402
protocol = p
398403
}
399404

400-
eps, zoneTarget = state.GetEndpointsByService(
401-
dataclientZone,
402-
ns,
403-
svcName,
404-
protocol,
405-
servicePort,
406-
)
407-
ic.logger.Debugf("Found %d endpoints for %s: %v", len(eps), svcName, err)
405+
if state.enableEndpointSlices {
406+
epSlices = state.GetEndpointSlicesByService(ns, svcName, protocol, servicePort)
407+
for _, ep := range epSlices {
408+
eps = append(eps, ep.Address)
409+
}
410+
} else {
411+
eps = state.GetEndpointsByService(
412+
ns,
413+
svcName,
414+
protocol,
415+
servicePort,
416+
)
417+
ic.logger.Debugf("Found %d endpoints for %s: %v", len(eps), svcName, err)
418+
}
408419
}
409420

410421
if len(eps) == 0 {
@@ -429,10 +440,10 @@ func (ing *ingress) convertDefaultBackendV1(
429440
LBAlgorithm: getLoadBalancerAlgorithm(i.Metadata, ing.defaultLoadBalancerAlgorithm),
430441
}
431442

432-
if zoneTarget {
443+
if state.enableEndpointSlices {
433444
var lbeps []*eskip.LBEndpoint
434-
for _, ep := range eps {
435-
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep, Zone: dataclientZone})
445+
for _, ep := range epSlices {
446+
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep.Address, Zone: ep.Zone})
436447
}
437448
r.LBEndpoints = lbeps
438449
} else {

dataclients/kubernetes/routegroup.go

Lines changed: 23 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -188,16 +188,26 @@ func applyServiceBackend(ctx *routeGroupContext, backend *definitions.SkipperBac
188188
return targetPortNotFound(backend.ServiceName, backend.ServicePort)
189189
}
190190

191-
dataclientZone := ctx.zone
192-
193-
eps, zoneTarget := ctx.state.GetEndpointsByTarget(
194-
dataclientZone,
195-
namespaceString(ctx.routeGroup.Metadata.Namespace),
196-
s.Meta.Name,
197-
"TCP",
198-
protocol,
199-
targetPort,
200-
)
191+
var eps []string
192+
var epSlices []skipperEndpoint
193+
if ctx.state.enableEndpointSlices {
194+
epSlices = ctx.state.GetEndpointSlicesByTarget(namespaceString(ctx.routeGroup.Metadata.Namespace),
195+
s.Meta.Name,
196+
"TCP",
197+
protocol,
198+
targetPort)
199+
for _, epSlice := range epSlices {
200+
eps = append(eps, epSlice.Address)
201+
}
202+
} else {
203+
eps = ctx.state.GetEndpointsByTarget(
204+
namespaceString(ctx.routeGroup.Metadata.Namespace),
205+
s.Meta.Name,
206+
"TCP",
207+
protocol,
208+
targetPort,
209+
)
210+
}
201211

202212
if len(eps) == 0 {
203213
ctx.logger.Tracef("Target endpoints not found, shuntroute for %s:%d", backend.ServiceName, backend.ServicePort)
@@ -213,9 +223,9 @@ func applyServiceBackend(ctx *routeGroupContext, backend *definitions.SkipperBac
213223
}
214224

215225
r.BackendType = eskip.LBBackend
216-
if zoneTarget {
217-
for _, ep := range eps {
218-
r.LBEndpoints = append(r.LBEndpoints, &eskip.LBEndpoint{Address: ep, Zone: dataclientZone})
226+
if ctx.state.enableEndpointSlices {
227+
for _, ep := range epSlices {
228+
r.LBEndpoints = append(r.LBEndpoints, &eskip.LBEndpoint{Address: ep.Address, Zone: ep.Zone})
219229
}
220230
} else {
221231
r.LBEndpoints = eskip.NewLBEndpoints(eps)

0 commit comments

Comments
 (0)