Skip to content

Commit 7cd729f

Browse files
authored
KAFKA-21074: Kafka protocol fault proxy clients fixtures (#23438)
A lightweight Kafka wire-protocol fault-injection proxy for integration tests. Sit it in front of an EmbeddedKafkaCluster, point any client’s bootstrap.servers at it, and it can inject error codes, drop connections, or blackhole a client — decoding/encoding with Kafka’s own protocol classes so it is correct across every wire version (including flexible/tagged-field ones), no hand-rolled byte offsets. Placed in `clients/src/testFixtures` so it is reusable from any module. This is internal test infrastructure only so it does not require a KIP. Fault primitives, each armed with a fluent, deterministic trigger: `injectError(apiKey, Errors)` — stamp an error code onto the matching response, re-serialized at the correct wire version. Supported APIs: `END_TXN`, `INIT_PRODUCER_ID`, `ADD_OFFSETS_TO_TXN`, `TXN_OFFSET_COMMIT`, `PRODUCE`, `FETCH` (extensible by registering a per-API setter). `disconnectOn(apiKey)` — drop the connection when the matching response would return (models the EOS "commit gap"). Works on any API. `delayOn(apiKey, Duration)` — hold back the matching response to model a slow broker, isolated to that one connection. `blackholeClient(clientIdSubstring)` — drop all requests from a targeted client before the broker sees them, so the broker evicts it by session timeout (a reversible, ungraceful one-node partition). Shared trigger DSL (`Occurrence`): `once()`, `onCall`, `times`, `everyTime()`, `withProbability(p)`. The deterministic triggers are safe for assertions; probability is chaos-mode only. The proxy never closes sockets unless a `disconnectOn(...)` rule fires, so it is not itself a source of flakiness. Rules can be armed/disarmed live from the test thread and each returns a handle exposing match/fire counts. Reviewers: Chia-Ping Tsai <chia7712@gmail.com>, Lucas Brutschy <lbrutschy@confluent.io>
1 parent 53abf0e commit 7cd729f

6 files changed

Lines changed: 1039 additions & 0 deletions

File tree

Lines changed: 201 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,201 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
package org.apache.kafka.test.faultproxy;
18+
19+
import org.apache.kafka.common.protocol.ApiKeys;
20+
import org.apache.kafka.common.protocol.Errors;
21+
22+
import java.util.concurrent.atomic.AtomicInteger;
23+
import java.util.function.IntPredicate;
24+
25+
/**
26+
* A single fault to apply to responses of one {@link ApiKeys}. Built through the fluent DSL on
27+
* {@link KafkaProtocolFaultProxy} ({@code injectError(...)} / {@code disconnectOn(...)} then a trigger),
28+
* and returned to the caller as a handle so a test can inspect {@link #timesTriggered()} or
29+
* {@link #remove()} it.
30+
*
31+
* <p>A rule counts each response offered to it ("matches"); the {@link IntPredicate} trigger decides, from
32+
* the 1-based match count, whether the fault fires on that match — e.g. {@code n -> n == 1} is {@code once()},
33+
* {@code n -> n <= 3} is {@code times(3)}. Responses are offered in registration order and consumed by the
34+
* first rule that fires, so a later rule is not offered a response an earlier rule already acted on.
35+
*/
36+
public final class FaultRule {
37+
38+
enum Action { INJECT_ERROR, DISCONNECT, DELAY }
39+
40+
private final KafkaProtocolFaultProxy owner;
41+
private final ApiKeys apiKey;
42+
private final Action action;
43+
private final Errors error; // only for INJECT_ERROR
44+
private final long delayMillis; // only for DELAY
45+
private final IntPredicate trigger;
46+
private final String clientIdFilter; // null = any client; otherwise the request clientId must contain it
47+
private final String description;
48+
49+
private final AtomicInteger matches = new AtomicInteger(0);
50+
private final AtomicInteger triggered = new AtomicInteger(0);
51+
52+
FaultRule(final KafkaProtocolFaultProxy owner,
53+
final ApiKeys apiKey,
54+
final Action action,
55+
final Errors error,
56+
final long delayMillis,
57+
final IntPredicate trigger,
58+
final String clientIdFilter,
59+
final String description) {
60+
this.owner = owner;
61+
this.apiKey = apiKey;
62+
this.action = action;
63+
this.error = error;
64+
this.delayMillis = delayMillis;
65+
this.trigger = trigger;
66+
this.clientIdFilter = clientIdFilter;
67+
this.description = description;
68+
}
69+
70+
ApiKeys apiKey() {
71+
return apiKey;
72+
}
73+
74+
/**
75+
* Whether this rule applies to a request from {@code clientId}. Lets a rule be scoped to a specific
76+
* client — e.g. {@code forClient("restore")} to fault only the restore consumer's fetches, leaving the
77+
* main consumer untouched. A null filter matches any client.
78+
*/
79+
boolean matchesClient(final String clientId) {
80+
return clientIdFilter == null || (clientId != null && clientId.contains(clientIdFilter));
81+
}
82+
83+
Action action() {
84+
return action;
85+
}
86+
87+
Errors error() {
88+
return error;
89+
}
90+
91+
long delayMillis() {
92+
return delayMillis;
93+
}
94+
95+
/** Called by the proxy for each response of this rule's API; returns true if the fault should fire. */
96+
boolean shouldFire() {
97+
final int n = matches.incrementAndGet();
98+
if (trigger.test(n)) {
99+
triggered.incrementAndGet();
100+
return true;
101+
}
102+
return false;
103+
}
104+
105+
/** How many times this fault actually fired — for test assertions. */
106+
public int timesTriggered() {
107+
return triggered.get();
108+
}
109+
110+
/** How many responses of this API the rule has observed (whether or not it fired). */
111+
public int timesMatched() {
112+
return matches.get();
113+
}
114+
115+
/** Deregister this rule from the proxy. */
116+
public void remove() {
117+
owner.removeFault(this);
118+
}
119+
120+
@Override
121+
public String toString() {
122+
return "FaultRule(" + description + ", triggered=" + triggered.get() + "/" + matches.get() + ")";
123+
}
124+
125+
// ------------------------------------------------------------------
126+
// Fluent trigger step: returned by injectError(...)/disconnectOn(...).
127+
// Each terminal registers the built rule with the proxy and returns the handle.
128+
// ------------------------------------------------------------------
129+
public static final class Builder {
130+
private final KafkaProtocolFaultProxy owner;
131+
private final ApiKeys apiKey;
132+
private final Action action;
133+
private final Errors error;
134+
private final long delayMillis;
135+
private String clientIdFilter; // null = any client
136+
137+
Builder(final KafkaProtocolFaultProxy owner, final ApiKeys apiKey, final Action action,
138+
final Errors error, final long delayMillis) {
139+
this.owner = owner;
140+
this.apiKey = apiKey;
141+
this.action = action;
142+
this.error = error;
143+
this.delayMillis = delayMillis;
144+
}
145+
146+
/**
147+
* Scope this fault to requests whose clientId contains {@code substring} — e.g.
148+
* {@code forClient("restore")} to fault only the restore consumer (leaving the main consumer and
149+
* producer untouched). Without this, the rule matches any client.
150+
*/
151+
public Builder forClient(final String substring) {
152+
this.clientIdFilter = substring;
153+
return this;
154+
}
155+
156+
private FaultRule register(final IntPredicate trigger, final String triggerDesc) {
157+
final String verb;
158+
switch (action) {
159+
case DISCONNECT:
160+
verb = "disconnect";
161+
break;
162+
case DELAY:
163+
verb = "delay " + delayMillis + "ms";
164+
break;
165+
default:
166+
verb = "inject " + error;
167+
break;
168+
}
169+
final String scope = clientIdFilter == null ? "" : " client~" + clientIdFilter;
170+
final FaultRule rule = new FaultRule(owner, apiKey, action, error, delayMillis, trigger, clientIdFilter,
171+
verb + " on " + apiKey + scope + " [" + triggerDesc + "]");
172+
owner.addFault(rule);
173+
return rule;
174+
}
175+
176+
/** Fire on the first matching response only. */
177+
public FaultRule once() {
178+
return register(Occurrence.once(), "once");
179+
}
180+
181+
/** Fire on exactly the {@code n}-th matching response (1-based). */
182+
public FaultRule onCall(final int n) {
183+
return register(Occurrence.call(n), "call #" + n);
184+
}
185+
186+
/** Fire on the first {@code n} matching responses. */
187+
public FaultRule times(final int n) {
188+
return register(Occurrence.first(n), "first " + n);
189+
}
190+
191+
/** Fire on every matching response until removed. */
192+
public FaultRule everyTime() {
193+
return register(Occurrence.always(), "every time");
194+
}
195+
196+
/** Fire on each matching response with the given probability (chaos mode; non-deterministic). */
197+
public FaultRule withProbability(final double p) {
198+
return register(Occurrence.withProbability(p), "p=" + p);
199+
}
200+
}
201+
}

0 commit comments

Comments
 (0)