Skip to content

Commit fa5b342

Browse files
authored
fix(broker): reconcile ingress contracts idempotently (#797)
Signed-off-by: kahirokunn <okinakahiro@gmail.com>
1 parent ac21ea7 commit fa5b342

5 files changed

Lines changed: 149 additions & 11 deletions

File tree

pkg/broker/contract/manager.go

Lines changed: 16 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,13 @@ func (m *Manager) GetContract(ctx context.Context) (*Contract, error) {
6161

6262
// UpdateBroker adds or updates a broker in the contract ConfigMap
6363
func (m *Manager) UpdateBroker(ctx context.Context, broker BrokerContract) error {
64+
_, err := m.UpdateBrokerIfChanged(ctx, broker)
65+
return err
66+
}
67+
68+
// UpdateBrokerIfChanged adds or updates a broker and reports whether it wrote
69+
// the contract ConfigMap.
70+
func (m *Manager) UpdateBrokerIfChanged(ctx context.Context, broker BrokerContract) (bool, error) {
6471
m.mu.Lock()
6572
defer m.mu.Unlock()
6673

@@ -74,7 +81,7 @@ func (m *Manager) UpdateBroker(ctx context.Context, broker BrokerContract) error
7481

7582
data, err := SerializeContract(contract)
7683
if err != nil {
77-
return err
84+
return false, err
7885
}
7986

8087
cm = &corev1.ConfigMap{
@@ -94,22 +101,24 @@ func (m *Manager) UpdateBroker(ctx context.Context, broker BrokerContract) error
94101
}
95102

96103
_, err = m.client.CoreV1().ConfigMaps(m.namespace).Create(ctx, cm, metav1.CreateOptions{})
97-
return err
104+
return err == nil, err
98105
}
99-
return err
106+
return false, err
100107
}
101108

102109
// Update existing ConfigMap
103110
contract, err := ParseContract(cm)
104111
if err != nil {
105-
return err
112+
return false, err
106113
}
107114

108-
contract.SetBroker(broker)
115+
if !contract.SetBrokerIfChanged(broker) {
116+
return false, nil
117+
}
109118

110119
data, err := SerializeContract(contract)
111120
if err != nil {
112-
return err
121+
return false, err
113122
}
114123

115124
if cm.Data == nil {
@@ -123,7 +132,7 @@ func (m *Manager) UpdateBroker(ctx context.Context, broker BrokerContract) error
123132
cm.Annotations[ContractGenerationAnnotation] = fmt.Sprintf("%d", contract.Generation)
124133

125134
_, err = m.client.CoreV1().ConfigMaps(m.namespace).Update(ctx, cm, metav1.UpdateOptions{})
126-
return err
135+
return err == nil, err
127136
}
128137

129138
// DeleteBroker removes a broker from the contract ConfigMap

pkg/broker/contract/manager_test.go

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -274,6 +274,88 @@ func TestManager_UpdateBroker_AnnotationUpdated(t *testing.T) {
274274
}
275275
}
276276

277+
func TestManager_UpdateBroker_SteadyStateDoesNotWrite(t *testing.T) {
278+
existingCM := contractConfigMap(&Contract{Brokers: make(map[string]BrokerContract)})
279+
client := fake.NewClientset(existingCM)
280+
manager := &Manager{
281+
client: client,
282+
lister: newFakeConfigMapLister(),
283+
namespace: testSystemNamespace,
284+
}
285+
broker := BrokerContract{
286+
UID: "uid-1", Namespace: "ns", Name: "br", StreamName: "stream-1",
287+
PublishSubject: "broker.ns.br", Path: "/ns/br", Generation: 4,
288+
}
289+
countWrites := func() int {
290+
writes := 0
291+
for _, action := range client.Actions() {
292+
if action.GetResource().Resource == "configmaps" &&
293+
(action.GetVerb() == "create" || action.GetVerb() == "update") {
294+
writes++
295+
}
296+
}
297+
return writes
298+
}
299+
getContract := func() *Contract {
300+
t.Helper()
301+
cm, err := client.CoreV1().ConfigMaps(testSystemNamespace).Get(
302+
context.Background(), ConfigMapName, metav1.GetOptions{},
303+
)
304+
if err != nil {
305+
t.Fatal(err)
306+
}
307+
contract, err := ParseContract(cm)
308+
if err != nil {
309+
t.Fatal(err)
310+
}
311+
return contract
312+
}
313+
314+
wrote, err := manager.UpdateBrokerIfChanged(context.Background(), broker)
315+
if err != nil {
316+
t.Fatal(err)
317+
}
318+
if !wrote {
319+
t.Fatal("initial UpdateBrokerIfChanged() changed = false, want true")
320+
}
321+
if got := getContract().Generation; got != 1 {
322+
t.Fatalf("Generation after initial UpdateBroker() = %d, want 1", got)
323+
}
324+
if got := countWrites(); got != 1 {
325+
t.Fatalf("ConfigMap writes after initial UpdateBroker() = %d, want 1", got)
326+
}
327+
328+
wrote, err = manager.UpdateBrokerIfChanged(context.Background(), broker)
329+
if err != nil {
330+
t.Fatal(err)
331+
}
332+
if wrote {
333+
t.Fatal("identical UpdateBrokerIfChanged() changed = true, want false")
334+
}
335+
if got := getContract().Generation; got != 1 {
336+
t.Fatalf("Generation after identical UpdateBroker() = %d, want unchanged 1", got)
337+
}
338+
if got := countWrites(); got != 1 {
339+
t.Fatalf("ConfigMap writes after identical UpdateBroker() = %d, want unchanged 1", got)
340+
}
341+
342+
changedContract := broker
343+
changedContract.PublishSubject = "broker.ns.br.changed"
344+
wrote, err = manager.UpdateBrokerIfChanged(context.Background(), changedContract)
345+
if err != nil {
346+
t.Fatal(err)
347+
}
348+
if !wrote {
349+
t.Fatal("changed UpdateBrokerIfChanged() changed = false, want true")
350+
}
351+
if got := getContract().Generation; got != 2 {
352+
t.Fatalf("Generation after changed UpdateBroker() = %d, want 2", got)
353+
}
354+
if got := countWrites(); got != 2 {
355+
t.Fatalf("ConfigMap writes after changed UpdateBroker() = %d, want 2", got)
356+
}
357+
}
358+
277359
// --- DeleteBroker tests ---
278360

279361
func TestManager_DeleteBroker_ConfigMapNotFound(t *testing.T) {

pkg/broker/contract/types.go

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -88,13 +88,24 @@ func (c *Contract) GetBrokerByPath(path string) (BrokerContract, bool) {
8888
return BrokerContract{}, false
8989
}
9090

91-
// SetBroker adds or updates a broker in the contract
91+
// SetBroker adds or updates a broker in the contract.
9292
func (c *Contract) SetBroker(broker BrokerContract) {
93+
c.SetBrokerIfChanged(broker)
94+
}
95+
96+
// SetBrokerIfChanged adds or updates a broker and reports whether the
97+
// serialized contract changed.
98+
func (c *Contract) SetBrokerIfChanged(broker BrokerContract) bool {
9399
if c.Brokers == nil {
94100
c.Brokers = make(map[string]BrokerContract)
95101
}
96-
c.Brokers[BrokerKey(broker.Namespace, broker.Name)] = broker
102+
key := BrokerKey(broker.Namespace, broker.Name)
103+
if existing, ok := c.Brokers[key]; ok && existing == broker {
104+
return false
105+
}
106+
c.Brokers[key] = broker
97107
c.Generation++
108+
return true
98109
}
99110

100111
// DeleteBroker removes a broker from the contract

pkg/broker/contract/types_test.go

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -129,6 +129,39 @@ func TestContract_SetBroker(t *testing.T) {
129129
t.Errorf("StreamName = %q, want %q", got, "new")
130130
}
131131
})
132+
133+
t.Run("identical broker is a steady-state no-op", func(t *testing.T) {
134+
broker := BrokerContract{
135+
UID: "uid-1", Namespace: "ns", Name: "br", StreamName: "stream-1",
136+
PublishSubject: "broker.ns.br", Path: "/ns/br", Generation: 4,
137+
}
138+
contract := &Contract{Brokers: make(map[string]BrokerContract)}
139+
140+
if changed := contract.SetBrokerIfChanged(broker); !changed {
141+
t.Fatal("initial SetBrokerIfChanged() changed = false, want true")
142+
}
143+
if contract.Generation != 1 {
144+
t.Fatalf("Generation after initial SetBroker() = %d, want 1", contract.Generation)
145+
}
146+
if changed := contract.SetBrokerIfChanged(broker); changed {
147+
t.Fatal("identical SetBrokerIfChanged() changed = true, want false")
148+
}
149+
if contract.Generation != 1 {
150+
t.Fatalf("Generation after identical SetBroker() = %d, want unchanged 1", contract.Generation)
151+
}
152+
153+
changed := broker
154+
changed.StreamName = "stream-2"
155+
if changedContract := contract.SetBrokerIfChanged(changed); !changedContract {
156+
t.Fatal("changed SetBrokerIfChanged() changed = false, want true")
157+
}
158+
if contract.Generation != 2 {
159+
t.Fatalf("Generation after changed SetBroker() = %d, want 2", contract.Generation)
160+
}
161+
if got := contract.Brokers["ns/br"].StreamName; got != "stream-2" {
162+
t.Fatalf("stored StreamName = %q, want %q", got, "stream-2")
163+
}
164+
})
132165
}
133166

134167
func TestContract_DeleteBroker(t *testing.T) {

pkg/broker/controller/reconciler.go

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -139,14 +139,17 @@ func (r *Reconciler) ReconcileKind(ctx context.Context, b *eventingv1.Broker) pk
139139
}
140140

141141
// Step 4: Reconcile ingress service
142-
if err := r.contractManager.UpdateBroker(ctx, brokerContract); err != nil {
142+
contractUpdated, err := r.contractManager.UpdateBrokerIfChanged(ctx, brokerContract)
143+
if err != nil {
143144
logger.Errorw("Failed to update contract", zap.Error(err))
144145
controller.GetEventRecorder(ctx).Event(b, corev1.EventTypeWarning, ReasonContractFailed, err.Error())
145146
b.Status.MarkIngressFailed("ContractUpdateFailed", "Failed to update contract ConfigMap: %v", err)
146147
return fmt.Errorf("failed to update contract: %w", err)
147148
}
148149

149-
controller.GetEventRecorder(ctx).Event(b, corev1.EventTypeNormal, ReasonContractUpdated, "Contract updated")
150+
if contractUpdated {
151+
controller.GetEventRecorder(ctx).Event(b, corev1.EventTypeNormal, ReasonContractUpdated, "Contract updated")
152+
}
150153

151154
// Step 4: Check shared ingress deployment readiness
152155
if err := r.propagateIngressAvailability(ctx, b); err != nil {

0 commit comments

Comments
 (0)