Skip to content

Commit b2efe51

Browse files
committed
kubernetes: use AWS physical zone IDs for locality
Signed-off-by: Martin Baillie <martin@baillie.id>
1 parent 0988daa commit b2efe51

8 files changed

Lines changed: 320 additions & 6 deletions

File tree

internal/provider/kubernetes/controller.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -626,6 +626,8 @@ func (r *gatewayAPIReconciler) Reconcile(ctx context.Context, _ reconcile.Reques
626626
}
627627
}
628628

629+
r.applyAWSZoneIDs(gwcResources)
630+
629631
// Sort before storing to:
630632
// 1. ensure identical resources are not retranslated
631633
// and updates are avoided by the watchable layer
@@ -2709,10 +2711,10 @@ func (r *gatewayAPIReconciler) watchResources(ctx context.Context, mgr manager.M
27092711
}
27102712
}
27112713

2712-
// Watch Node CRUDs to update Gateway Address exposed by Service of type NodePort.
2714+
// Watch Nodes for locality changes and Gateway addresses exposed by NodePort Services.
27132715
// Node creation/deletion and ExternalIP updates would require update in the Gateway
27142716
nPredicates := []predicate.TypedPredicate[*corev1.Node]{
2715-
predicate.TypedGenerationChangedPredicate[*corev1.Node]{},
2717+
r.nodePredicate(),
27162718
predicate.NewTypedPredicateFuncs(func(node *corev1.Node) bool {
27172719
return r.handleNode(node)
27182720
}),
Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,69 @@
1+
// Copyright Envoy Gateway Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
// The full text of the Apache license is available in the LICENSE file at
4+
// the root of the repo.
5+
6+
package kubernetes
7+
8+
import (
9+
corev1 "k8s.io/api/core/v1"
10+
discoveryv1 "k8s.io/api/discovery/v1"
11+
"sigs.k8s.io/controller-runtime/pkg/event"
12+
"sigs.k8s.io/controller-runtime/pkg/predicate"
13+
14+
"github.com/envoyproxy/gateway/internal/gatewayapi/resource"
15+
)
16+
17+
const (
18+
awsZoneIDLabel = "topology.k8s.aws/zone-id"
19+
endpointSliceControllerName = "endpointslice-controller.k8s.io"
20+
endpointSliceMirroringControllerName = "endpointslicemirroring-controller.k8s.io"
21+
)
22+
23+
// applyAWSZoneIDs resolves local endpoints without changing cached EndpointSlices.
24+
func (r *gatewayAPIReconciler) applyAWSZoneIDs(resources resource.ControllerResources) {
25+
zoneIDs := make(map[string]string)
26+
for _, res := range resources {
27+
for i, slice := range res.EndpointSlices {
28+
// Slices from other controllers can refer to nodes in another cluster.
29+
managedBy := slice.Labels[discoveryv1.LabelManagedBy]
30+
if managedBy != endpointSliceControllerName &&
31+
managedBy != endpointSliceMirroringControllerName {
32+
continue
33+
}
34+
var updated *discoveryv1.EndpointSlice
35+
for j, endpoint := range slice.Endpoints {
36+
if endpoint.NodeName == nil || *endpoint.NodeName == "" {
37+
continue
38+
}
39+
zoneID, found := zoneIDs[*endpoint.NodeName]
40+
if !found {
41+
zoneID = r.store.nodeZoneID(*endpoint.NodeName)
42+
zoneIDs[*endpoint.NodeName] = zoneID
43+
}
44+
if zoneID == "" || (endpoint.Zone != nil && *endpoint.Zone == zoneID) {
45+
continue
46+
}
47+
if updated == nil {
48+
updated = slice.DeepCopy()
49+
}
50+
updated.Endpoints[j].Zone = new(zoneID)
51+
}
52+
if updated != nil {
53+
res.EndpointSlices[i] = updated
54+
}
55+
}
56+
}
57+
}
58+
59+
func (r *gatewayAPIReconciler) nodePredicate() predicate.TypedPredicate[*corev1.Node] {
60+
return predicate.Or(
61+
predicate.TypedGenerationChangedPredicate[*corev1.Node]{},
62+
predicate.TypedFuncs[*corev1.Node]{
63+
UpdateFunc: func(e event.TypedUpdateEvent[*corev1.Node]) bool {
64+
return e.ObjectOld != nil && e.ObjectNew != nil &&
65+
e.ObjectOld.Labels[awsZoneIDLabel] != e.ObjectNew.Labels[awsZoneIDLabel]
66+
},
67+
},
68+
)
69+
}
Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
1+
// Copyright Envoy Gateway Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
// The full text of the Apache license is available in the LICENSE file at
4+
// the root of the repo.
5+
6+
package kubernetes
7+
8+
import (
9+
"fmt"
10+
"os"
11+
"testing"
12+
13+
endpointv3 "github.com/envoyproxy/go-control-plane/envoy/config/endpoint/v3"
14+
resourcev3 "github.com/envoyproxy/go-control-plane/pkg/resource/v3"
15+
"github.com/stretchr/testify/require"
16+
corev1 "k8s.io/api/core/v1"
17+
discoveryv1 "k8s.io/api/discovery/v1"
18+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
19+
"k8s.io/apimachinery/pkg/types"
20+
"sigs.k8s.io/controller-runtime/pkg/client"
21+
"sigs.k8s.io/controller-runtime/pkg/event"
22+
mcsapiv1a1 "sigs.k8s.io/mcs-api/pkg/apis/v1alpha1"
23+
24+
egv1a1 "github.com/envoyproxy/gateway/api/v1alpha1"
25+
"github.com/envoyproxy/gateway/internal/envoygateway/config"
26+
"github.com/envoyproxy/gateway/internal/gatewayapi"
27+
"github.com/envoyproxy/gateway/internal/gatewayapi/resource"
28+
"github.com/envoyproxy/gateway/internal/infrastructure/kubernetes/proxy"
29+
"github.com/envoyproxy/gateway/internal/message"
30+
"github.com/envoyproxy/gateway/internal/provider/kubernetes/test"
31+
xdstranslator "github.com/envoyproxy/gateway/internal/xds/translator"
32+
)
33+
34+
func TestApplyAWSZoneIDs(t *testing.T) {
35+
local := &discoveryv1.EndpointSlice{
36+
ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{discoveryv1.LabelManagedBy: endpointSliceControllerName}},
37+
Endpoints: []discoveryv1.Endpoint{
38+
{NodeName: new("node"), Zone: new("us-east-1d")},
39+
{NodeName: new("node")},
40+
{NodeName: new("no-id"), Zone: new("us-east-1b")},
41+
{NodeName: new("empty-id"), Zone: new("us-east-1b")},
42+
{NodeName: new("missing"), Zone: new("us-east-1c")},
43+
{Zone: new("explicit")},
44+
},
45+
}
46+
local.Labels[mcsapiv1a1.LabelServiceName] = "also-exported"
47+
imported := local.DeepCopy()
48+
imported.Labels[mcsapiv1a1.LabelServiceName] = "remote"
49+
imported.Labels[discoveryv1.LabelManagedBy] = "mcs-controller"
50+
custom := local.DeepCopy()
51+
custom.Labels[discoveryv1.LabelManagedBy] = "custom-controller"
52+
mirrored := local.DeepCopy()
53+
mirrored.Labels[discoveryv1.LabelManagedBy] = endpointSliceMirroringControllerName
54+
original := local.DeepCopy()
55+
resources := resource.ControllerResources{
56+
{EndpointSlices: []*discoveryv1.EndpointSlice{local, imported, custom, mirrored}},
57+
{EndpointSlices: []*discoveryv1.EndpointSlice{local}},
58+
}
59+
store := newProviderStore()
60+
for _, node := range []*corev1.Node{
61+
{ObjectMeta: metav1.ObjectMeta{Name: "node", Labels: map[string]string{awsZoneIDLabel: "use1-az6"}}},
62+
{ObjectMeta: metav1.ObjectMeta{Name: "no-id"}},
63+
{ObjectMeta: metav1.ObjectMeta{Name: "empty-id", Labels: map[string]string{awsZoneIDLabel: ""}}},
64+
} {
65+
store.addNode(node)
66+
}
67+
r := &gatewayAPIReconciler{store: store}
68+
r.applyAWSZoneIDs(resources)
69+
for _, res := range resources {
70+
require.Equal(t, "use1-az6", *res.EndpointSlices[0].Endpoints[0].Zone)
71+
require.Equal(t, "use1-az6", *res.EndpointSlices[0].Endpoints[1].Zone)
72+
require.Equal(t, original.Endpoints[2:], res.EndpointSlices[0].Endpoints[2:])
73+
}
74+
require.Equal(t, original, local)
75+
require.Equal(t, "use1-az6", *resources[0].EndpointSlices[3].Endpoints[0].Zone)
76+
require.Same(t, imported, resources[0].EndpointSlices[1])
77+
require.Same(t, custom, resources[0].EndpointSlices[2])
78+
require.Equal(t, original.Endpoints, imported.Endpoints)
79+
require.Equal(t, original.Endpoints, custom.Endpoints)
80+
}
81+
82+
func TestAWSZoneIDNodeUpdates(t *testing.T) {
83+
r := &gatewayAPIReconciler{}
84+
old := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node"}}
85+
updated := old.DeepCopy()
86+
updated.Labels = map[string]string{awsZoneIDLabel: "use1-az6"}
87+
for _, nodes := range [][2]*corev1.Node{{old, updated}, {updated, old}} {
88+
require.True(t, r.nodePredicate().Update(event.TypedUpdateEvent[*corev1.Node]{ObjectOld: nodes[0], ObjectNew: nodes[1]}))
89+
}
90+
require.False(t, r.nodePredicate().Update(event.TypedUpdateEvent[*corev1.Node]{ObjectOld: updated, ObjectNew: updated}))
91+
updated.Generation++
92+
require.True(t, r.nodePredicate().Update(event.TypedUpdateEvent[*corev1.Node]{ObjectOld: old, ObjectNew: updated}))
93+
}
94+
95+
func TestAWSZoneIDLocalityXDS(t *testing.T) {
96+
cfg, err := config.New(os.Stdout, os.Stderr)
97+
require.NoError(t, err)
98+
cfg.EnvoyGateway.Provider = &egv1a1.EnvoyGatewayProvider{Type: egv1a1.ProviderTypeCustom}
99+
updates := new(message.ProviderResources)
100+
r, err := NewOfflineGatewayAPIController(t.Context(), cfg, nil, updates)
101+
require.NoError(t, err)
102+
defer updates.Close()
103+
gateway := test.GetGateway(types.NamespacedName{Namespace: "default", Name: "gateway"}, "eg", 80)
104+
fleet := test.GetService(types.NamespacedName{
105+
Namespace: cfg.ControllerNamespace, Name: proxy.ExpectedResourceHashedName("default/gateway"),
106+
}, gatewayapi.OwnerLabels(gateway, false), map[string]int32{"dummy": 8080})
107+
backend := test.GetService(types.NamespacedName{Namespace: "default", Name: "backend"}, nil, map[string]int32{"dummy": 8080})
108+
backend.Spec.TrafficDistribution = new("PreferClose")
109+
fleetSlice := test.GetEndpointSlice(types.NamespacedName{Namespace: fleet.Namespace, Name: "fleet"}, fleet.Name, false)
110+
backendSlice := test.GetEndpointSlice(types.NamespacedName{Namespace: backend.Namespace, Name: "backend"}, backend.Name, false)
111+
for i, slice := range []*discoveryv1.EndpointSlice{fleetSlice, backendSlice} {
112+
slice.AddressType = discoveryv1.AddressTypeIPv4
113+
slice.Labels[discoveryv1.LabelManagedBy] = endpointSliceControllerName
114+
slice.Endpoints[0].NodeName = new("node")
115+
slice.Endpoints[0].Zone = new("us-east-1d")
116+
slice.Endpoints[0].Addresses = []string{fmt.Sprintf("192.0.2.%d", i+1)}
117+
}
118+
remoteSlice := backendSlice.DeepCopy()
119+
remoteSlice.Name = "remote"
120+
remoteSlice.Labels[discoveryv1.LabelManagedBy] = "remote-controller"
121+
remoteSlice.Endpoints[0].Zone = new("use1-az2")
122+
remoteSlice.Endpoints[0].Addresses = []string{"192.0.2.3"}
123+
node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node", Labels: map[string]string{
124+
corev1.LabelTopologyZone: "us-east-1d", awsZoneIDLabel: "use1-az6",
125+
}}}
126+
for _, obj := range []client.Object{
127+
&corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "default"}},
128+
&corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: cfg.ControllerNamespace}},
129+
test.GetGatewayClass("eg", egv1a1.GatewayControllerName, nil), gateway, fleet, backend,
130+
fleetSlice, backendSlice, remoteSlice, node,
131+
test.GetHTTPRoute(types.NamespacedName{Namespace: "default", Name: "route"}, gateway.Name,
132+
test.GetServiceBackendRef(types.NamespacedName{Name: backend.Name}, 8080), ""),
133+
} {
134+
require.NoError(t, r.Client.Create(t.Context(), obj))
135+
}
136+
r.store.addNode(node)
137+
require.NoError(t, r.Reconcile(t.Context()))
138+
published, ok := updates.GatewayAPIResources.Load(cfg.EnvoyGateway.Gateway.ControllerName)
139+
require.True(t, ok)
140+
require.Len(t, *published.Resources, 1)
141+
gt := &gatewayapi.Translator{
142+
GatewayControllerName: egv1a1.GatewayControllerName, GatewayClassName: "eg",
143+
ControllerNamespace: cfg.ControllerNamespace, Logger: cfg.Logger,
144+
}
145+
translated, err := gt.Translate(t.Context(), (*published.Resources)[0])
146+
require.NoError(t, err)
147+
require.Len(t, translated.XdsIR, 1)
148+
zones := map[string]string{}
149+
for _, xdsIR := range translated.XdsIR {
150+
xt := &xdstranslator.Translator{Logger: cfg.Logger}
151+
table, err := xt.Translate(t.Context(), xdsIR)
152+
require.NoError(t, err)
153+
for _, entry := range table.XdsResources[resourcev3.EndpointType] {
154+
for _, locality := range entry.(*endpointv3.ClusterLoadAssignment).Endpoints {
155+
for _, endpoint := range locality.LbEndpoints {
156+
zones[endpoint.GetEndpoint().Address.GetSocketAddress().Address] = locality.Locality.GetZone()
157+
}
158+
}
159+
}
160+
}
161+
require.Equal(t, map[string]string{
162+
"192.0.2.1": "use1-az6", "192.0.2.2": "use1-az6", "192.0.2.3": "use1-az2",
163+
}, zones)
164+
var stored discoveryv1.EndpointSlice
165+
require.NoError(t, r.Client.Get(t.Context(), client.ObjectKeyFromObject(backendSlice), &stored))
166+
require.Equal(t, "us-east-1d", *stored.Endpoints[0].Zone)
167+
}

internal/provider/kubernetes/store.go

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ import (
1717
type nodeDetails struct {
1818
name string
1919
addresses status.NodeAddresses
20+
zoneID string
2021
}
2122

2223
// kubernetesProviderStore holds cached information for the kubernetes provider.
@@ -35,7 +36,7 @@ func newProviderStore() *kubernetesProviderStore {
3536
}
3637

3738
func (p *kubernetesProviderStore) addNode(n *corev1.Node) {
38-
details := nodeDetails{name: n.Name}
39+
details := nodeDetails{name: n.Name, zoneID: n.Labels[awsZoneIDLabel]}
3940

4041
var internalIPs, externalIPs status.NodeAddresses
4142
for _, addr := range n.Status.Addresses {
@@ -74,6 +75,12 @@ func (p *kubernetesProviderStore) removeNode(n *corev1.Node) {
7475
delete(p.nodes, n.Name)
7576
}
7677

78+
func (p *kubernetesProviderStore) nodeZoneID(name string) string {
79+
p.mu.Lock()
80+
defer p.mu.Unlock()
81+
return p.nodes[name].zoneID
82+
}
83+
7784
func (p *kubernetesProviderStore) listNodeAddresses() status.NodeAddresses {
7885
p.mu.Lock()
7986
defer p.mu.Unlock()

internal/provider/kubernetes/topology_injector.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,12 @@ func (m *ProxyTopologyInjector) Handle(ctx context.Context, req admission.Reques
8383
}
8484
logger = logger.WithValues("node", node)
8585

86-
if zone, ok := node.Labels[corev1.LabelTopologyZone]; ok {
86+
zone, hasZone := node.Labels[corev1.LabelTopologyZone]
87+
if node.Labels[awsZoneIDLabel] != "" {
88+
zone = node.Labels[awsZoneIDLabel]
89+
hasZone = true
90+
}
91+
if hasZone {
8792
if binding.Annotations == nil {
8893
binding.Annotations = map[string]string{}
8994
}

internal/provider/kubernetes/topology_injector_test.go

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,17 @@ func TestProxyTopologyInjector_Handle(t *testing.T) {
4545
},
4646
}
4747

48+
awsNode := defaultNode.DeepCopy()
49+
awsNode.Labels = map[string]string{
50+
corev1.LabelTopologyZone: "us-east-1d",
51+
"topology.k8s.aws/zone-id": "use1-az6",
52+
}
53+
54+
awsOnlyNode := awsNode.DeepCopy()
55+
delete(awsOnlyNode.Labels, corev1.LabelTopologyZone)
56+
emptyAWSNode := awsNode.DeepCopy()
57+
emptyAWSNode.Labels[awsZoneIDLabel] = ""
58+
4859
cases := []struct {
4960
caseName string
5061
obj client.Object
@@ -71,6 +82,48 @@ func TestProxyTopologyInjector_Handle(t *testing.T) {
7182
},
7283
}},
7384
},
85+
{
86+
caseName: "AWS zone ID",
87+
obj: &corev1.Binding{
88+
ObjectMeta: metav1.ObjectMeta{Name: defaultPod.Name, Namespace: defaultPod.Namespace},
89+
Target: corev1.ObjectReference{Name: awsNode.Name},
90+
},
91+
node: awsNode,
92+
pod: defaultPod,
93+
expectedPatchResp: []jsonpatch.JsonPatchOperation{{
94+
Operation: "add",
95+
Path: "/metadata/annotations",
96+
Value: map[string]interface{}{corev1.LabelTopologyZone: "\"use1-az6\""},
97+
}},
98+
},
99+
{
100+
caseName: "AWS zone ID without standard label",
101+
obj: &corev1.Binding{
102+
ObjectMeta: metav1.ObjectMeta{Name: defaultPod.Name, Namespace: defaultPod.Namespace},
103+
Target: corev1.ObjectReference{Name: awsNode.Name},
104+
},
105+
node: awsOnlyNode,
106+
pod: defaultPod,
107+
expectedPatchResp: []jsonpatch.JsonPatchOperation{{
108+
Operation: "add",
109+
Path: "/metadata/annotations",
110+
Value: map[string]interface{}{corev1.LabelTopologyZone: "\"use1-az6\""},
111+
}},
112+
},
113+
{
114+
caseName: "empty AWS zone ID",
115+
obj: &corev1.Binding{
116+
ObjectMeta: metav1.ObjectMeta{Name: defaultPod.Name, Namespace: defaultPod.Namespace},
117+
Target: corev1.ObjectReference{Name: awsNode.Name},
118+
},
119+
node: emptyAWSNode,
120+
pod: defaultPod,
121+
expectedPatchResp: []jsonpatch.JsonPatchOperation{{
122+
Operation: "add",
123+
Path: "/metadata/annotations",
124+
Value: map[string]interface{}{corev1.LabelTopologyZone: "\"us-east-1d\""},
125+
}},
126+
},
74127
{
75128
caseName: "empty target",
76129
obj: &corev1.Binding{
Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
Use physical AWS zone IDs for proxy and local Kubernetes-managed endpoint locality when the node
2+
label is available. Update Backend zones, `weightedZones`, and imported or custom EndpointSlices to
3+
use IDs.

0 commit comments

Comments
 (0)