Skip to content

Commit 37de7a9

Browse files
committed
feat: remove no-longer-desired edges on topology change (cleanup gating)
Signed-off-by: Shreesha001 <shettyshreesha552@gmail.com>
1 parent 170645c commit 37de7a9

4 files changed

Lines changed: 129 additions & 14 deletions

File tree

‎service/topology_resolver.go‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,3 +67,32 @@ func ResolveTopologyEdges(clusters []string, topology *controllerv1alpha1.Topolo
6767
}
6868
return edges
6969
}
70+
71+
// TopologyEdgeSet answers direction-insensitive membership questions about a
72+
// set of desired edges: the two WorkerSliceGateway objects of a pair (the
73+
// server side and the client side) belong to the same logical edge.
74+
type TopologyEdgeSet struct {
75+
members map[[2]string]bool
76+
}
77+
78+
// NewTopologyEdgeSet builds a TopologyEdgeSet from resolved edges.
79+
func NewTopologyEdgeSet(edges []TopologyEdge) TopologyEdgeSet {
80+
members := make(map[[2]string]bool, len(edges))
81+
for _, edge := range edges {
82+
members[edgeKey(edge.ServerCluster, edge.ClientCluster)] = true
83+
}
84+
return TopologyEdgeSet{members: members}
85+
}
86+
87+
// Contains reports whether the given cluster pair, in either order, is a
88+
// desired edge.
89+
func (s TopologyEdgeSet) Contains(clusterA, clusterB string) bool {
90+
return s.members[edgeKey(clusterA, clusterB)]
91+
}
92+
93+
func edgeKey(clusterA, clusterB string) [2]string {
94+
if clusterA > clusterB {
95+
clusterA, clusterB = clusterB, clusterA
96+
}
97+
return [2]string{clusterA, clusterB}
98+
}

‎service/topology_resolver_test.go‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,3 +78,22 @@ func TestResolveTopologyEdges(t *testing.T) {
7878
})
7979
}
8080
}
81+
82+
func TestTopologyEdgeSetContains(t *testing.T) {
83+
// hub-and-spoke edges: hub=worker-1, spokes worker-2/worker-3
84+
set := NewTopologyEdgeSet(ResolveTopologyEdges(
85+
[]string{"worker-1", "worker-2", "worker-3"},
86+
&controllerv1alpha1.TopologySpec{Mode: controllerv1alpha1.TopologyModeHubAndSpoke, Hubs: []string{"worker-1"}},
87+
))
88+
// desired hub<->spoke edges, both directions
89+
if !set.Contains("worker-1", "worker-2") || !set.Contains("worker-2", "worker-1") {
90+
t.Fatal("expected worker-1<->worker-2 to be a desired edge (either direction)")
91+
}
92+
if !set.Contains("worker-1", "worker-3") {
93+
t.Fatal("expected worker-1<->worker-3 to be a desired edge")
94+
}
95+
// spoke<->spoke is NOT desired
96+
if set.Contains("worker-2", "worker-3") || set.Contains("worker-3", "worker-2") {
97+
t.Fatal("did not expect worker-2<->worker-3 (spoke-to-spoke) to be a desired edge")
98+
}
99+
}

‎service/worker_slice_gateway_service.go‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -352,15 +352,15 @@ func (s *WorkerSliceGatewayService) CreateMinimumWorkerSliceGateways(ctx context
352352
sliceSubnet string, clusterCidr string, sliceGwSvcTypeMap map[string]*controllerv1alpha1.SliceGatewayServiceType,
353353
topology *controllerv1alpha1.TopologySpec) (ctrl.Result, error) {
354354

355-
err := s.cleanupObsoleteGateways(ctx, namespace, label, clusterNames, clusterMap)
355+
desiredEdges := ResolveTopologyEdges(clusterNames, topology)
356+
err := s.cleanupObsoleteGateways(ctx, namespace, label, clusterNames, clusterMap, NewTopologyEdgeSet(desiredEdges))
356357
if err != nil {
357358
return ctrl.Result{}, err
358359
}
359360
if len(clusterNames) < 2 {
360361
return ctrl.Result{}, nil
361362
}
362363

363-
desiredEdges := ResolveTopologyEdges(clusterNames, topology)
364364
_, err = s.createMinimumGatewaysIfNotExists(ctx, sliceName, desiredEdges, namespace, label, clusterMap, sliceSubnet, clusterCidr, sliceGwSvcTypeMap)
365365
if err != nil {
366366
return ctrl.Result{}, err
@@ -381,7 +381,7 @@ func (s *WorkerSliceGatewayService) ListWorkerSliceGateways(ctx context.Context,
381381

382382
// cleanupObsoleteGateways is a function delete outdated gateways
383383
func (s *WorkerSliceGatewayService) cleanupObsoleteGateways(ctx context.Context, namespace string, ownerLabel map[string]string,
384-
clusters []string, clusterMap map[string]int) error {
384+
clusters []string, clusterMap map[string]int, desiredEdges TopologyEdgeSet) error {
385385

386386
gateways, err := s.ListWorkerSliceGateways(ctx, ownerLabel, namespace)
387387
if err != nil {
@@ -408,7 +408,10 @@ func (s *WorkerSliceGatewayService) cleanupObsoleteGateways(ctx context.Context,
408408
clusterSource := gateway.Spec.LocalGatewayConfig.ClusterName
409409
clusterDestination := gateway.Spec.RemoteGatewayConfig.ClusterName
410410
gatewayExpectedNumber := s.calculateGatewayNumber(clusterMap[clusterSource], clusterMap[clusterDestination])
411-
if !clusterExistMap[clusterSource] || !clusterExistMap[clusterDestination] || gatewayExpectedNumber != gateway.Spec.GatewayNumber {
411+
// Delete a gateway when either cluster left the slice, its gateway number
412+
// changed, or its edge is no longer in the desired topology (e.g. a
413+
// spoke<->spoke link after a FullMesh->HubAndSpoke change).
414+
if !clusterExistMap[clusterSource] || !clusterExistMap[clusterDestination] || gatewayExpectedNumber != gateway.Spec.GatewayNumber || !desiredEdges.Contains(clusterSource, clusterDestination) {
412415
err = util.DeleteResource(ctx, &gateway)
413416
if err != nil {
414417
//Register an event for worker slice gateway deletion failure

‎service/worker_slice_gateway_service_test.go‎

Lines changed: 74 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -57,16 +57,17 @@ func TestWorkerSliceGatewaySuite(t *testing.T) {
5757
}
5858

5959
var WorkerSliceGatewayTestbed = map[string]func(*testing.T){
60-
"TestWorkerSliceGatewayReconciliation_Success": testWorkerSliceGatewayReconciliationSuccess,
61-
"TestWorkerSliceGatewayReconciliation_IfSliceConfigNotFound": testWorkerSliceGatewayReconciliationIfSliceConfigNotFound,
62-
"TestWorkerSliceGatewayReconciliation_IfGatewayNotFound": testWorkerSliceGatewayReconciliationIfGatewayNotFound,
63-
"TestWorkerSliceGatewayReconciliation_Delete": testWorkerSliceGatewayReconciliationDelete,
64-
"TestWorkerSliceGatewayReconciliation_DeleteForcefully": testWorkerSliceGatewayReconciliationDeleteForcefully,
65-
"TestCreateMinimumWorkerSliceGateways_IfAlreadyExists": testCreateMinimumWorkerSliceGatewaysAlreadyExists,
66-
"TestCreateMinimumWorkerSliceGateways_HubAndSpokeSkipsSpokeToSpoke": testCreateMinimumWorkerSliceGatewaysHubAndSpokeSkipsSpokeToSpoke,
67-
"TestCreateMinimumWorkerSliceGateways_IfNotExists": testCreateMinimumWorkerSliceGatewaysNotExists,
68-
"TestDeleteWorkerSliceGatewaysByLabel_IfExists": testDeleteWorkerSliceGatewaysByLabelExists,
69-
"TestNodeIpReconciliationOfWorkerSliceGateways_IfExists": testNodeIpReconciliationOfWorkerSliceGatewaysExists,
60+
"TestWorkerSliceGatewayReconciliation_Success": testWorkerSliceGatewayReconciliationSuccess,
61+
"TestWorkerSliceGatewayReconciliation_IfSliceConfigNotFound": testWorkerSliceGatewayReconciliationIfSliceConfigNotFound,
62+
"TestWorkerSliceGatewayReconciliation_IfGatewayNotFound": testWorkerSliceGatewayReconciliationIfGatewayNotFound,
63+
"TestWorkerSliceGatewayReconciliation_Delete": testWorkerSliceGatewayReconciliationDelete,
64+
"TestWorkerSliceGatewayReconciliation_DeleteForcefully": testWorkerSliceGatewayReconciliationDeleteForcefully,
65+
"TestCreateMinimumWorkerSliceGateways_IfAlreadyExists": testCreateMinimumWorkerSliceGatewaysAlreadyExists,
66+
"TestCreateMinimumWorkerSliceGateways_HubAndSpokeSkipsSpokeToSpoke": testCreateMinimumWorkerSliceGatewaysHubAndSpokeSkipsSpokeToSpoke,
67+
"TestCreateMinimumWorkerSliceGateways_HubAndSpokeCleansUpSpokeToSpoke": testCreateMinimumWorkerSliceGatewaysHubAndSpokeCleansUpSpokeToSpoke,
68+
"TestCreateMinimumWorkerSliceGateways_IfNotExists": testCreateMinimumWorkerSliceGatewaysNotExists,
69+
"TestDeleteWorkerSliceGatewaysByLabel_IfExists": testDeleteWorkerSliceGatewaysByLabelExists,
70+
"TestNodeIpReconciliationOfWorkerSliceGateways_IfExists": testNodeIpReconciliationOfWorkerSliceGatewaysExists,
7071
}
7172

7273
func testWorkerSliceGatewayReconciliationSuccess(t *testing.T) {
@@ -369,6 +370,69 @@ func testCreateMinimumWorkerSliceGatewaysHubAndSpokeSkipsSpokeToSpoke(t *testing
369370
mMock.AssertExpectations(t)
370371
}
371372

373+
// testCreateMinimumWorkerSliceGatewaysHubAndSpokeCleansUpSpokeToSpoke verifies
374+
// that when a slice already has a spoke<->spoke gateway pair (e.g. left over from
375+
// a FullMesh->HubAndSpoke change), cleanup deletes it purely because its edge is
376+
// no longer in the desired hub-and-spoke set. The stale pair is given the CORRECT
377+
// gateway number and both clusters are still slice members, so the ONLY reason it
378+
// gets removed is the topology edge check.
379+
func testCreateMinimumWorkerSliceGatewaysHubAndSpokeCleansUpSpokeToSpoke(t *testing.T) {
380+
_, _, _, workerSliceGatewayService, requestObj, clientMock, _, ctx, mMock := setupWorkerSliceGatewayTest("slice_gateway", "namespace")
381+
label := map[string]string{}
382+
clusterNames := []string{"cluster-1", "cluster-2", "cluster-3"}
383+
clusterMap := map[string]int{
384+
"cluster-1": 1,
385+
"cluster-2": 2,
386+
"cluster-3": 3,
387+
}
388+
topology := &controllerv1alpha1.TopologySpec{
389+
Mode: controllerv1alpha1.TopologyModeHubAndSpoke,
390+
Hubs: []string{"cluster-1"},
391+
}
392+
mMock.On("WithProject", mock.AnythingOfType("string")).Return(&metrics.MetricRecorder{}).Once()
393+
// existing spoke<->spoke pair (cluster-2 <-> cluster-3) with the CORRECT gateway
394+
// number; both clusters are still members, so it survives the membership and
395+
// number checks and is removed only because its edge is not desired.
396+
spokeToSpokeNumber := ((3-1)*(3-2))/2 + 2 // calculateGatewayNumber(2, 3) = 3
397+
pairWorkerSliceGateway := &workerv1alpha1.WorkerSliceGatewayList{}
398+
clientMock.On("List", ctx, pairWorkerSliceGateway, mock.Anything, client.InNamespace(requestObj.Namespace)).Return(nil).Run(func(args mock.Arguments) {
399+
arg := args.Get(1).(*workerv1alpha1.WorkerSliceGatewayList)
400+
arg.Items = []workerv1alpha1.WorkerSliceGateway{
401+
{
402+
Spec: workerv1alpha1.WorkerSliceGatewaySpec{
403+
LocalGatewayConfig: workerv1alpha1.SliceGatewayConfig{ClusterName: "cluster-2"},
404+
RemoteGatewayConfig: workerv1alpha1.SliceGatewayConfig{ClusterName: "cluster-3"},
405+
GatewayNumber: spokeToSpokeNumber,
406+
},
407+
},
408+
{
409+
Spec: workerv1alpha1.WorkerSliceGatewaySpec{
410+
LocalGatewayConfig: workerv1alpha1.SliceGatewayConfig{ClusterName: "cluster-3"},
411+
RemoteGatewayConfig: workerv1alpha1.SliceGatewayConfig{ClusterName: "cluster-2"},
412+
GatewayNumber: spokeToSpokeNumber,
413+
},
414+
},
415+
}
416+
}).Once()
417+
clientMock.On("Delete", ctx, mock.Anything).Return(nil).Twice()
418+
clientMock.On("Create", ctx, mock.AnythingOfType("*v1.Event")).Return(nil).Once()
419+
mMock.On("RecordCounterMetric", mock.Anything, mock.Anything).Return().Once()
420+
clientMock.On("Update", ctx, mock.AnythingOfType("*v1.Event")).Return(nil).Once()
421+
mMock.On("RecordCounterMetric", mock.Anything, mock.Anything).Return().Once()
422+
// create pass for the two desired hub<->spoke edges: clusters fetched, gateways
423+
// already exist -> nothing created.
424+
cluster := &controllerv1alpha1.Cluster{}
425+
clientMock.On("Get", ctx, mock.AnythingOfType("types.NamespacedName"), cluster).Return(nil).Times(3)
426+
gateway := &workerv1alpha1.WorkerSliceGateway{}
427+
clientMock.On("Get", ctx, mock.AnythingOfType("types.NamespacedName"), gateway).Return(nil).Times(4)
428+
429+
result, err := workerSliceGatewayService.CreateMinimumWorkerSliceGateways(ctx, "red", clusterNames, requestObj.Namespace, label, clusterMap, "10.10.10.10/16", "/16", nil, topology)
430+
require.Equal(t, ctrl.Result{}, result)
431+
require.Nil(t, err)
432+
clientMock.AssertExpectations(t)
433+
mMock.AssertExpectations(t)
434+
}
435+
372436
func testCreateMinimumWorkerSliceGatewaysNotExists(t *testing.T) {
373437
_, _, jobMock, workerSliceGatewayService, requestObj, clientMock, _, ctx, mMock := setupWorkerSliceGatewayTest("slice_gateway", "namespace")
374438
label := map[string]string{

0 commit comments

Comments
 (0)