Skip to content

Commit b1f9037

Browse files
authored
[kubernetes/leaderelection] Fix race condition (#38694)
1 parent abb5cfc commit b1f9037

3 files changed

Lines changed: 187 additions & 3 deletions

File tree

pkg/util/kubernetes/apiserver/leaderelection/leaderelection_engine.go

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,10 +106,12 @@ func (le *LeaderEngine) createLeaderTokenIfNotExists() error {
106106
Namespace: le.LeaderNamespace,
107107
},
108108
}, metav1.CreateOptions{})
109-
if err != nil && !errors.IsConflict(err) {
109+
if err != nil && !errors.IsConflict(err) && !errors.IsAlreadyExists(err) {
110110
return err
111111
}
112112
}
113+
114+
return nil
113115
}
114116
_, err := le.coreClient.ConfigMaps(le.LeaderNamespace).Get(context.TODO(), le.LeaseName, metav1.GetOptions{})
115117
if err != nil {
@@ -125,11 +127,11 @@ func (le *LeaderEngine) createLeaderTokenIfNotExists() error {
125127
Name: le.LeaseName,
126128
},
127129
}, metav1.CreateOptions{})
128-
if err != nil && !errors.IsConflict(err) {
130+
if err != nil && !errors.IsConflict(err) && !errors.IsAlreadyExists(err) {
129131
return err
130132
}
131133
}
132-
return err
134+
return nil
133135
}
134136

135137
// newElection creates an election.
Lines changed: 166 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,166 @@
1+
// Unless explicitly stated otherwise all files in this repository are licensed
2+
// under the Apache License Version 2.0.
3+
// This product includes software developed at Datadog (https://www.datadoghq.com/).
4+
// Copyright 2016-present Datadog, Inc.
5+
6+
//go:build kubeapiserver
7+
8+
package leaderelection
9+
10+
import (
11+
"context"
12+
"testing"
13+
14+
"github.com/stretchr/testify/assert"
15+
"github.com/stretchr/testify/require"
16+
coordinationv1 "k8s.io/api/coordination/v1"
17+
corev1 "k8s.io/api/core/v1"
18+
"k8s.io/apimachinery/pkg/api/errors"
19+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
20+
"k8s.io/apimachinery/pkg/runtime"
21+
"k8s.io/client-go/kubernetes/fake"
22+
k8stesting "k8s.io/client-go/testing"
23+
rl "k8s.io/client-go/tools/leaderelection/resourcelock"
24+
25+
cmLock "github.com/DataDog/datadog-agent/internal/third_party/client-go/tools/leaderelection/resourcelock"
26+
)
27+
28+
func TestCreateLeaderTokenIfNotExists(t *testing.T) {
29+
tokenNamespace := "default"
30+
tokenName := "test-lease"
31+
32+
tests := []struct {
33+
name string
34+
lockType string
35+
setupFunc func(*fake.Clientset) error
36+
expectsError bool
37+
verifyFunc func(*testing.T, *fake.Clientset)
38+
}{
39+
{
40+
name: "Leases - lease does not exist",
41+
lockType: rl.LeasesResourceLock,
42+
expectsError: false,
43+
verifyFunc: func(t *testing.T, client *fake.Clientset) {
44+
_, err := client.CoordinationV1().Leases(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
45+
require.NoError(t, err)
46+
47+
// Regression test: check that it does not create a ConfigMap too
48+
_, err = client.CoreV1().ConfigMaps(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
49+
require.True(t, errors.IsNotFound(err))
50+
},
51+
},
52+
{
53+
name: "ConfigMap - configmap does not exist",
54+
lockType: cmLock.ConfigMapsResourceLock,
55+
expectsError: false,
56+
verifyFunc: func(t *testing.T, client *fake.Clientset) {
57+
_, err := client.CoreV1().ConfigMaps(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
58+
require.NoError(t, err)
59+
60+
_, err = client.CoordinationV1().Leases(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
61+
require.True(t, errors.IsNotFound(err))
62+
},
63+
},
64+
{
65+
name: "Leases - lease already exists",
66+
lockType: rl.LeasesResourceLock,
67+
setupFunc: func(client *fake.Clientset) error {
68+
lease := &coordinationv1.Lease{
69+
ObjectMeta: metav1.ObjectMeta{
70+
Name: tokenName,
71+
Namespace: tokenNamespace,
72+
},
73+
Spec: coordinationv1.LeaseSpec{},
74+
}
75+
_, err := client.CoordinationV1().Leases(tokenNamespace).Create(context.TODO(), lease, metav1.CreateOptions{})
76+
return err
77+
},
78+
expectsError: false,
79+
verifyFunc: func(t *testing.T, client *fake.Clientset) {
80+
_, err := client.CoordinationV1().Leases(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
81+
require.NoError(t, err)
82+
},
83+
},
84+
{
85+
name: "ConfigMap - configmap already exists",
86+
lockType: cmLock.ConfigMapsResourceLock,
87+
setupFunc: func(client *fake.Clientset) error {
88+
configMap := &corev1.ConfigMap{
89+
ObjectMeta: metav1.ObjectMeta{
90+
Name: tokenName,
91+
Namespace: tokenNamespace,
92+
},
93+
}
94+
_, err := client.CoreV1().ConfigMaps(tokenNamespace).Create(context.TODO(), configMap, metav1.CreateOptions{})
95+
return err
96+
},
97+
expectsError: false,
98+
verifyFunc: func(t *testing.T, client *fake.Clientset) {
99+
_, err := client.CoreV1().ConfigMaps(tokenNamespace).Get(context.TODO(), tokenName, metav1.GetOptions{})
100+
require.NoError(t, err)
101+
},
102+
},
103+
{
104+
name: "Leases - lease created between check and create",
105+
lockType: rl.LeasesResourceLock,
106+
setupFunc: func(client *fake.Clientset) error {
107+
// Mock the Create method to return an "AlreadyExists" error. This
108+
// simulates another DCA creating the resource between the check and
109+
// the creation operations.
110+
client.PrependReactor("create", "leases", func(_ k8stesting.Action) (bool, runtime.Object, error) {
111+
return true, nil, errors.NewAlreadyExists(coordinationv1.Resource("leases"), tokenName)
112+
})
113+
114+
return nil
115+
},
116+
expectsError: false, // Should handle the race condition and not return error
117+
},
118+
{
119+
name: "ConfigMap - configmap created between check and create",
120+
lockType: cmLock.ConfigMapsResourceLock,
121+
setupFunc: func(client *fake.Clientset) error {
122+
// Mock the Create method to return an "AlreadyExists" error. This
123+
// simulates another DCA creating the resource between the check and
124+
// the creation operations.
125+
client.PrependReactor("create", "configmaps", func(_ k8stesting.Action) (bool, runtime.Object, error) {
126+
return true, nil, errors.NewAlreadyExists(corev1.Resource("configmaps"), tokenName)
127+
})
128+
129+
return nil
130+
},
131+
expectsError: false, // Should handle the race condition and not return error
132+
},
133+
}
134+
135+
for _, test := range tests {
136+
t.Run(test.name, func(t *testing.T) {
137+
client := fake.NewClientset()
138+
139+
le := &LeaderEngine{
140+
LeaseName: tokenName,
141+
LeaderNamespace: tokenNamespace,
142+
coreClient: client.CoreV1(),
143+
coordClient: client.CoordinationV1(),
144+
lockType: test.lockType,
145+
}
146+
147+
if test.setupFunc != nil {
148+
err := test.setupFunc(client)
149+
require.NoError(t, err)
150+
}
151+
152+
err := le.createLeaderTokenIfNotExists()
153+
154+
if test.expectsError {
155+
assert.Error(t, err)
156+
return
157+
}
158+
159+
assert.NoError(t, err)
160+
161+
if test.verifyFunc != nil {
162+
test.verifyFunc(t, client)
163+
}
164+
})
165+
}
166+
}
Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
# Each section from every release note are combined when the
2+
# CHANGELOG.rst is rendered. So the text needs to be worded so that
3+
# it does not depend on any information only available in another
4+
# section. This may mean repeating some details, but each section
5+
# must be readable independently of the other.
6+
#
7+
# Each section note must be formatted as reStructuredText.
8+
---
9+
fixes:
10+
- |
11+
Fixed an issue that caused some Datadog Cluster Agent replicas to restart
12+
during the first install.
13+
- |
14+
When using leases for leader election with the
15+
``leader_election_default_resource`` option, the Datadog Cluster Agent no
16+
longer creates an unused config map.

0 commit comments

Comments
 (0)