Skip to content

Commit 5f162ef

Browse files
committed
use confimap to track gateway infra
Signed-off-by: Huabing (Robin) Zhao <zhaohuabing@gmail.com>
1 parent 8821341 commit 5f162ef

5 files changed

Lines changed: 435 additions & 1 deletion

File tree

internal/infrastructure/kubernetes/infra.go

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,10 @@ type Infra struct {
5959
// Client wrap k8s client.
6060
Client *InfraClient
6161

62+
// APIReader performs uncached reads/writes for narrow correctness-sensitive
63+
// tracking state in the controller namespace. Defaults to the main client in tests.
64+
APIReader client.Client
65+
6266
logger logging.Logger
6367

6468
// errors is the notifier used to send async errors to the main control loop.
@@ -74,6 +78,7 @@ func NewInfra(cli client.Client, cfg *config.Server, errors message.RunnerErrorN
7478
DNSDomain: cfg.DNSDomain,
7579
EnvoyGateway: cfg.EnvoyGateway,
7680
Client: New(cli),
81+
APIReader: cli,
7782
logger: cfg.Logger.WithName(string(egv1a1.LogComponentInfrastructureRunner)),
7883
errors: errors,
7984
}
@@ -85,6 +90,31 @@ func (i *Infra) Close() error { return nil }
8590
// createOrUpdate creates a ServiceAccount/ConfigMap/Deployment/Service in the kube api server based on the
8691
// provided ResourceRender, if it doesn't exist and updates it if it does.
8792
func (i *Infra) createOrUpdate(ctx context.Context, r ResourceRender) error {
93+
// Track the last applied proxy resource names in a controller-owned ConfigMap
94+
// so renamed resources can be deleted directly by name on the next reconcile.
95+
// This avoids relying solely on DeleteAllExcept's cached List view to
96+
// discover stale custom-named resources, while keeping that path as a
97+
// fallback for older untracked leftovers.
98+
var (
99+
previousRefs []proxy.ManagedProxyResourceRef
100+
currentRefs []proxy.ManagedProxyResourceRef
101+
trackingName string
102+
err error
103+
)
104+
if proxyRender, ok := r.(*proxy.ResourceRender); ok {
105+
trackingName = proxyRender.TrackingConfigMapName()
106+
if len(trackingName) > 0 {
107+
previousRefs, err = i.loadManagedProxyRefs(ctx, trackingName)
108+
if err != nil {
109+
return fmt.Errorf("failed to load tracked proxy resource refs for %s: %w", trackingName, err)
110+
}
111+
}
112+
currentRefs, err = proxyRender.ManagedProxyResourceRefs()
113+
if err != nil {
114+
return fmt.Errorf("failed to compute managed proxy resource refs for %s/%s: %w", proxyRender.Namespace(), proxyRender.Name(), err)
115+
}
116+
}
117+
88118
if err := i.createOrUpdateServiceAccount(ctx, r); err != nil {
89119
return fmt.Errorf("failed to create or update serviceaccount %s/%s: %w", r.Namespace(), r.Name(), err)
90120
}
@@ -113,6 +143,15 @@ func (i *Infra) createOrUpdate(ctx context.Context, r ResourceRender) error {
113143
return fmt.Errorf("failed to create or update pdb %s/%s: %w", r.Namespace(), r.Name(), err)
114144
}
115145

146+
if len(trackingName) > 0 {
147+
if err := i.deleteTrackedProxyResources(ctx, previousRefs, currentRefs); err != nil {
148+
return fmt.Errorf("failed to delete previously tracked proxy resources for %s: %w", trackingName, err)
149+
}
150+
if err := i.persistManagedProxyRefs(ctx, trackingName, currentRefs); err != nil {
151+
return fmt.Errorf("failed to persist tracked proxy resources for %s: %w", trackingName, err)
152+
}
153+
}
154+
116155
return nil
117156
}
118157

@@ -146,5 +185,13 @@ func (i *Infra) delete(ctx context.Context, r ResourceRender) error {
146185
return fmt.Errorf("failed to delete pdb %s/%s: %w", r.Namespace(), r.Name(), err)
147186
}
148187

188+
if proxyRender, ok := r.(*proxy.ResourceRender); ok {
189+
if trackingName := proxyRender.TrackingConfigMapName(); len(trackingName) > 0 {
190+
if err := i.deleteTrackedProxyRefState(ctx, trackingName); err != nil {
191+
return fmt.Errorf("failed to delete tracked proxy ref state %s: %w", trackingName, err)
192+
}
193+
}
194+
}
195+
149196
return nil
150197
}

internal/infrastructure/kubernetes/proxy/resource_provider.go

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ package proxy
77

88
import (
99
"context"
10+
"crypto/sha256"
1011
"fmt"
1112
"strconv"
1213
"time"
@@ -77,6 +78,12 @@ type ResourceRender struct {
7778
ownerReferenceUID map[string]types.UID
7879
}
7980

81+
type ManagedProxyResourceRef struct {
82+
Kind string `json:"kind"`
83+
Namespace string `json:"namespace"`
84+
Name string `json:"name"`
85+
}
86+
8087
// KubernetesInfraProvider provide information for initializing the proxy resource render.
8188
type KubernetesInfraProvider interface {
8289
GetControllerNamespace() string
@@ -135,6 +142,98 @@ func (r *ResourceRender) LabelSelector() labels.Selector {
135142
return labels.SelectorFromSet(r.stableSelector().MatchLabels)
136143
}
137144

145+
func (r *ResourceRender) OwningGatewayNN() *types.NamespacedName {
146+
labels := r.infra.GetProxyMetadata().Labels
147+
name := labels[gatewayapi.OwningGatewayNameLabel]
148+
namespace := labels[gatewayapi.OwningGatewayNamespaceLabel]
149+
if len(name) == 0 || len(namespace) == 0 {
150+
return nil
151+
}
152+
153+
return &types.NamespacedName{Namespace: namespace, Name: name}
154+
}
155+
156+
// TrackingOwner returns the logical owner identity used to persist proxy infra
157+
// tracking state. In normal mode infra is rendered per Gateway even if the
158+
// effective EnvoyProxy config comes from a GatewayClass, so tracked cleanup must
159+
// stay isolated per Gateway. Merged gateways intentionally render one shared
160+
// infra set for the GatewayClass and only carry OwningGatewayClassLabel.
161+
func (r *ResourceRender) TrackingOwner() (kind string, nn types.NamespacedName, ok bool) {
162+
if gwNN := r.OwningGatewayNN(); gwNN != nil {
163+
return "Gateway", *gwNN, true
164+
}
165+
166+
gcName := r.infra.GetProxyMetadata().Labels[gatewayapi.OwningGatewayClassLabel]
167+
if len(gcName) == 0 {
168+
return "", types.NamespacedName{}, false
169+
}
170+
171+
return "GatewayClass", types.NamespacedName{Name: gcName}, true
172+
}
173+
174+
func (r *ResourceRender) TrackingConfigMapName() string {
175+
kind, nn, ok := r.TrackingOwner()
176+
if !ok {
177+
return ""
178+
}
179+
180+
key := kind + "/" + nn.Namespace + "/" + nn.Name
181+
sum := sha256.Sum256([]byte(key))
182+
return fmt.Sprintf("envoy-proxy-tracking-%x", sum[:8])
183+
}
184+
185+
func (r *ResourceRender) ManagedProxyResourceRefs() ([]ManagedProxyResourceRef, error) {
186+
refs := make([]ManagedProxyResourceRef, 0, 5)
187+
188+
sa, err := r.ServiceAccount()
189+
if err != nil {
190+
return nil, err
191+
}
192+
refs = append(refs, ManagedProxyResourceRef{Kind: "ServiceAccount", Namespace: sa.Namespace, Name: sa.Name})
193+
194+
svc, err := r.Service()
195+
if err != nil {
196+
return nil, err
197+
}
198+
if svc != nil {
199+
refs = append(refs, ManagedProxyResourceRef{Kind: "Service", Namespace: svc.Namespace, Name: svc.Name})
200+
}
201+
202+
deployment, err := r.Deployment()
203+
if err != nil {
204+
return nil, err
205+
}
206+
if deployment != nil {
207+
refs = append(refs, ManagedProxyResourceRef{Kind: "Deployment", Namespace: deployment.Namespace, Name: deployment.Name})
208+
}
209+
210+
daemonSet, err := r.DaemonSet()
211+
if err != nil {
212+
return nil, err
213+
}
214+
if daemonSet != nil {
215+
refs = append(refs, ManagedProxyResourceRef{Kind: "DaemonSet", Namespace: daemonSet.Namespace, Name: daemonSet.Name})
216+
}
217+
218+
hpa, err := r.HorizontalPodAutoscaler()
219+
if err != nil {
220+
return nil, err
221+
}
222+
if hpa != nil {
223+
refs = append(refs, ManagedProxyResourceRef{Kind: "HorizontalPodAutoscaler", Namespace: hpa.Namespace, Name: hpa.Name})
224+
}
225+
226+
pdb, err := r.PodDisruptionBudget()
227+
if err != nil {
228+
return nil, err
229+
}
230+
if pdb != nil {
231+
refs = append(refs, ManagedProxyResourceRef{Kind: "PodDisruptionBudget", Namespace: pdb.Namespace, Name: pdb.Name})
232+
}
233+
234+
return refs, nil
235+
}
236+
138237
func (r *ResourceRender) ownerReferences() []metav1.OwnerReference {
139238
var ownerReferences []metav1.OwnerReference
140239
if r.ownerReferenceUID != nil {

internal/infrastructure/kubernetes/proxy_infra_test.go

Lines changed: 129 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ package kubernetes
77

88
import (
99
"context"
10+
"encoding/json"
1011
"errors"
1112
"fmt"
1213
"os"
@@ -282,6 +283,90 @@ func TestDeleteProxyInfra(t *testing.T) {
282283
}
283284
}
284285

286+
func TestCreateProxyInfraDeletesTrackedServiceAccountsWhenListIsStale(t *testing.T) {
287+
listStaleInterceptor := interceptor.Funcs{
288+
Patch: interceptorFunc.Patch,
289+
List: func(ctx context.Context, clnt client.WithWatch, list client.ObjectList, opts ...client.ListOption) error {
290+
if err := clnt.List(ctx, list, opts...); err != nil {
291+
return err
292+
}
293+
saList, ok := list.(*corev1.ServiceAccountList)
294+
if ok && len(saList.Items) > 1 {
295+
saList.Items = saList.Items[:1]
296+
}
297+
return nil
298+
},
299+
}
300+
301+
cli := fakeclient.NewClientBuilder().
302+
WithScheme(envoygateway.GetScheme()).
303+
WithInterceptorFuncs(listStaleInterceptor).
304+
Build()
305+
kube := newTestInfraWithClient(t, cli)
306+
kube.EnvoyGateway.Provider = &egv1a1.EnvoyGatewayProvider{
307+
Type: egv1a1.ProviderTypeKubernetes,
308+
Kubernetes: &egv1a1.EnvoyGatewayKubernetesProvider{
309+
Deploy: &egv1a1.KubernetesDeployMode{
310+
Type: new(egv1a1.KubernetesDeployModeTypeGatewayNamespace),
311+
},
312+
},
313+
}
314+
315+
ctx := context.Background()
316+
require.NoError(t, setupOwnerReferenceResources(ctx, kube.Client))
317+
318+
gwNN := types.NamespacedName{Namespace: "ns1", Name: "gateway-1"}
319+
320+
defaultInfra := newProxyInfraForGateway(gwNN, "")
321+
require.NoError(t, kube.CreateOrUpdateProxyInfra(ctx, defaultInfra))
322+
assertTrackingConfigMap(t, ctx, kube, defaultInfra, []proxy.ManagedProxyResourceRef{
323+
{Kind: "ServiceAccount", Namespace: gwNN.Namespace, Name: gwNN.Name},
324+
{Kind: "Service", Namespace: gwNN.Namespace, Name: gwNN.Name},
325+
{Kind: "Deployment", Namespace: gwNN.Namespace, Name: gwNN.Name},
326+
})
327+
328+
customAInfra := newProxyInfraForGateway(gwNN, "custom-a-sa")
329+
require.NoError(t, kube.CreateOrUpdateProxyInfra(ctx, customAInfra))
330+
assertTrackingConfigMap(t, ctx, kube, customAInfra, []proxy.ManagedProxyResourceRef{
331+
{Kind: "ServiceAccount", Namespace: gwNN.Namespace, Name: "custom-a-sa"},
332+
{Kind: "Service", Namespace: gwNN.Namespace, Name: gwNN.Name},
333+
{Kind: "Deployment", Namespace: gwNN.Namespace, Name: gwNN.Name},
334+
})
335+
require.NoError(t, kube.Client.Get(ctx, client.ObjectKey{Namespace: gwNN.Namespace, Name: "custom-a-sa"}, &corev1.ServiceAccount{}))
336+
require.True(t, kerrors.IsNotFound(kube.Client.Get(ctx, gwNN, &corev1.ServiceAccount{})))
337+
338+
customBInfra := newProxyInfraForGateway(gwNN, "custom-b-sa")
339+
require.NoError(t, kube.CreateOrUpdateProxyInfra(ctx, customBInfra))
340+
assertTrackingConfigMap(t, ctx, kube, customBInfra, []proxy.ManagedProxyResourceRef{
341+
{Kind: "ServiceAccount", Namespace: gwNN.Namespace, Name: "custom-b-sa"},
342+
{Kind: "Service", Namespace: gwNN.Namespace, Name: gwNN.Name},
343+
{Kind: "Deployment", Namespace: gwNN.Namespace, Name: gwNN.Name},
344+
})
345+
require.NoError(t, kube.Client.Get(ctx, client.ObjectKey{Namespace: gwNN.Namespace, Name: "custom-b-sa"}, &corev1.ServiceAccount{}))
346+
require.True(t, kerrors.IsNotFound(kube.Client.Get(ctx, client.ObjectKey{Namespace: gwNN.Namespace, Name: "custom-a-sa"}, &corev1.ServiceAccount{})))
347+
}
348+
349+
func TestCreateProxyInfraTracksRefsInControllerNamespaceConfigMapForMergedGateways(t *testing.T) {
350+
kube := newTestInfra(t)
351+
ctx := context.Background()
352+
require.NoError(t, setupOwnerReferenceResources(ctx, kube.Client))
353+
354+
infra := ir.NewInfra()
355+
infra.Proxy.Metadata.Labels = proxy.EnvoyAppLabel()
356+
infra.Proxy.Metadata.Labels[gatewayapi.OwningGatewayClassLabel] = testGatewayClass
357+
infra.Proxy.Metadata.OwnerReference = &ir.ResourceMetadata{
358+
Kind: resource.KindGatewayClass,
359+
Name: testGatewayClass,
360+
}
361+
362+
require.NoError(t, kube.CreateOrUpdateProxyInfra(ctx, infra))
363+
assertTrackingConfigMap(t, ctx, kube, infra, []proxy.ManagedProxyResourceRef{
364+
{Kind: "ServiceAccount", Namespace: kube.ControllerNamespace, Name: expectedName(infra.Proxy, false)},
365+
{Kind: "Service", Namespace: kube.ControllerNamespace, Name: expectedName(infra.Proxy, false)},
366+
{Kind: "Deployment", Namespace: kube.ControllerNamespace, Name: expectedName(infra.Proxy, false)},
367+
})
368+
}
369+
285370
// This function uses setup creating Resources for OwnerReference.
286371
// When the default case, ProxyInfra Get OwnerReference from GatewayClass.
287372
// When enable GatewayNamespace mode, ProxyInfra Get OwnerReference from Gateway.
@@ -304,3 +389,47 @@ func setupOwnerReferenceResources(ctx context.Context, client *InfraClient) erro
304389
}
305390
return client.Create(ctx, gw)
306391
}
392+
393+
func newProxyInfraForGateway(gwNN types.NamespacedName, serviceAccountName string) *ir.Infra {
394+
infra := ir.NewInfra()
395+
infra.Proxy.Name = gwNN.Name
396+
infra.Proxy.Namespace = gwNN.Namespace
397+
infra.Proxy.Metadata.Labels = proxy.EnvoyAppLabel()
398+
infra.Proxy.Metadata.Labels[gatewayapi.OwningGatewayNamespaceLabel] = gwNN.Namespace
399+
infra.Proxy.Metadata.Labels[gatewayapi.OwningGatewayNameLabel] = gwNN.Name
400+
infra.Proxy.Metadata.OwnerReference = &ir.ResourceMetadata{
401+
Kind: resource.KindGateway,
402+
Name: gwNN.Name,
403+
}
404+
infra.Proxy.Config = &egv1a1.EnvoyProxy{
405+
Spec: egv1a1.EnvoyProxySpec{
406+
Provider: &egv1a1.EnvoyProxyProvider{
407+
Type: egv1a1.EnvoyProxyProviderTypeKubernetes,
408+
Kubernetes: egv1a1.DefaultEnvoyProxyKubeProvider(),
409+
},
410+
},
411+
}
412+
if serviceAccountName != "" {
413+
infra.Proxy.Config.Spec.Provider.Kubernetes.EnvoyServiceAccount = &egv1a1.KubernetesServiceAccountSpec{
414+
Name: &serviceAccountName,
415+
}
416+
}
417+
return infra
418+
}
419+
420+
func assertTrackingConfigMap(t *testing.T, ctx context.Context, kube *Infra, infra *ir.Infra, expected []proxy.ManagedProxyResourceRef) {
421+
t.Helper()
422+
423+
render, err := proxy.NewResourceRender(ctx, kube, infra)
424+
require.NoError(t, err)
425+
426+
cm := &corev1.ConfigMap{}
427+
require.NoError(t, kube.APIReader.Get(ctx, client.ObjectKey{
428+
Namespace: kube.ControllerNamespace,
429+
Name: render.TrackingConfigMapName(),
430+
}, cm))
431+
432+
var actual []proxy.ManagedProxyResourceRef
433+
require.NoError(t, json.Unmarshal([]byte(cm.Data[trackedProxyRefsDataKey]), &actual))
434+
require.Equal(t, expected, actual)
435+
}

0 commit comments

Comments
 (0)