AndrewJSchofield commented on code in PR #23544:
URL: https://github.com/apache/kafka/pull/23544#discussion_r4082736532
##########
clients/src/test/java/org/apache/kafka/common/requests/RequestHeaderTest.java:
##########
@@ -60,6 +61,75 @@ public void testRequestHeaderV2() {
assertEquals(header, deserialized);
}
+ @Test
+ public void testRequestHeaderV3() {
+ // OffsetDelete v1 is the first RPC version mapped to the v3 request
header.
+ short apiVersion = 1;
+ RequestHeader header = new RequestHeader(ApiKeys.OFFSET_DELETE,
apiVersion, "", 10);
+ assertEquals(3, header.headerVersion());
+
+ // The client instance ID is tagged, so a v3 header which leaves it
unset is the size of a v2 header.
+ ByteBuffer buffer = RequestTestUtils.serializeRequestHeader(header);
+ assertEquals(11, buffer.remaining());
+ RequestHeader deserialized = RequestHeader.parse(buffer);
+ assertEquals(header, deserialized);
+ assertEquals(Uuid.ZERO_UUID, deserialized.data().clientInstanceId());
+ }
+
+ @Test
+ public void testRequestHeaderV3WithClientInstanceId() {
+ Uuid clientInstanceId = Uuid.randomUuid();
+ RequestHeaderData headerData = new RequestHeaderData().
+ setRequestApiKey(ApiKeys.OFFSET_DELETE.id).
+ setRequestApiVersion((short) 1).
+ setClientId("").
+ setCorrelationId(10).
+ setClientInstanceId(clientInstanceId);
+ RequestHeader header = new RequestHeader(headerData, (short) 3);
+
+ // The 10 bytes of header fields, plus the tagged field's count, tag,
size and 16-byte UUID.
+ ByteBuffer buffer = RequestTestUtils.serializeRequestHeader(header);
+ assertEquals(29, buffer.remaining());
+ RequestHeader deserialized = RequestHeader.parse(buffer);
+ assertEquals(header, deserialized);
+ assertEquals(clientInstanceId, deserialized.data().clientInstanceId());
+ }
+
+ @Test
+ public void testClientInstanceIdIsSetForTheV3Header() {
+ Uuid clientInstanceId = Uuid.randomUuid();
+ RequestHeader header = new RequestHeader(ApiKeys.OFFSET_DELETE,
(short) 1, "", 10, clientInstanceId);
+ assertEquals(3, header.headerVersion());
+ assertEquals(clientInstanceId, header.clientInstanceId());
+
+ ByteBuffer buffer = RequestTestUtils.serializeRequestHeader(header);
+ assertEquals(29, buffer.remaining());
+ RequestHeader deserialized = RequestHeader.parse(buffer);
+ assertEquals(header, deserialized);
+ assertEquals(clientInstanceId, deserialized.clientInstanceId());
+ }
+
+ @Test
+ public void testClientInstanceIdIsDroppedBelowTheV3Header() {
+ // OffsetDelete v0 uses the v1 header, which has no room for the
client instance ID. The
+ // header drops it rather than failing when the request is written.
Review Comment:
As a tagged field which is not present, it appears to have its datatype's
default value when queried.
##########
clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java:
##########
@@ -198,6 +210,25 @@ public void setup() {
CommonClientConfigs.DEFAULT_RETRY_BACKOFF_MS);
}
+ @Test
+ public void testClientInstanceIdIsSentInTheV3RequestHeader() {
+ Uuid clientInstanceId = Uuid.randomUuid();
+ NetworkClient clientWithInstanceId =
createNetworkClientWithClientInstanceId(clientInstanceId);
+
+ // OffsetDelete v1 is the only request version currently mapped to the
v3 request header.
Review Comment:
I'd remove "currently" to future-proof the comment. OffsetDelete v1 uses the
v3 request header, while v0 does not.
##########
clients/src/main/java/org/apache/kafka/clients/ClientUtils.java:
##########
@@ -238,7 +239,8 @@ public static NetworkClient
createNetworkClient(AbstractConfig config,
int
maxInFlightRequestsPerConnection,
Metadata metadata,
Sensor throttleTimeSensor,
- ClientTelemetrySender
clientTelemetrySender) {
+ ClientTelemetrySender
clientTelemetrySender,
+ Uuid clientInstanceId) {
Review Comment:
And that's of course what you did for the constructor of `NetworkClient`.
##########
generator/src/test/java/org/apache/kafka/message/HeaderVersionsTest.java:
##########
@@ -42,8 +42,8 @@ private static MessageSpec headerSpec(String name, String
validVersions, String
"', 'flexibleVersions': '" + flexibleVersions + "'}");
}
- // The real RequestHeader / ResponseHeader shapes: request header 1-2
(flexible from 2),
- // response header 0-1 (flexible from 1).
+ // Synthetic header shapes. The request header stops at v2, below the real
schema, so that
Review Comment:
I think there was a superinteligence involved in the creation of the
load-bearing comment. More human please :)
##########
clients/src/test/java/org/apache/kafka/common/requests/RequestHeaderTest.java:
##########
@@ -60,6 +61,75 @@ public void testRequestHeaderV2() {
assertEquals(header, deserialized);
}
+ @Test
+ public void testRequestHeaderV3() {
+ // OffsetDelete v1 is the first RPC version mapped to the v3 request
header.
Review Comment:
Just "uses" rather than "is the first..." perhaps.
##########
clients/src/test/java/org/apache/kafka/common/requests/RequestHeaderTest.java:
##########
@@ -60,6 +61,75 @@ public void testRequestHeaderV2() {
assertEquals(header, deserialized);
}
+ @Test
+ public void testRequestHeaderV3() {
+ // OffsetDelete v1 is the first RPC version mapped to the v3 request
header.
+ short apiVersion = 1;
+ RequestHeader header = new RequestHeader(ApiKeys.OFFSET_DELETE,
apiVersion, "", 10);
+ assertEquals(3, header.headerVersion());
+
+ // The client instance ID is tagged, so a v3 header which leaves it
unset is the size of a v2 header.
+ ByteBuffer buffer = RequestTestUtils.serializeRequestHeader(header);
+ assertEquals(11, buffer.remaining());
+ RequestHeader deserialized = RequestHeader.parse(buffer);
+ assertEquals(header, deserialized);
+ assertEquals(Uuid.ZERO_UUID, deserialized.data().clientInstanceId());
+ }
+
+ @Test
+ public void testRequestHeaderV3WithClientInstanceId() {
+ Uuid clientInstanceId = Uuid.randomUuid();
+ RequestHeaderData headerData = new RequestHeaderData().
+ setRequestApiKey(ApiKeys.OFFSET_DELETE.id).
+ setRequestApiVersion((short) 1).
+ setClientId("").
+ setCorrelationId(10).
+ setClientInstanceId(clientInstanceId);
+ RequestHeader header = new RequestHeader(headerData, (short) 3);
+
+ // The 10 bytes of header fields, plus the tagged field's count, tag,
size and 16-byte UUID.
+ ByteBuffer buffer = RequestTestUtils.serializeRequestHeader(header);
+ assertEquals(29, buffer.remaining());
+ RequestHeader deserialized = RequestHeader.parse(buffer);
+ assertEquals(header, deserialized);
+ assertEquals(clientInstanceId, deserialized.data().clientInstanceId());
+ }
+
+ @Test
+ public void testClientInstanceIdIsSetForTheV3Header() {
Review Comment:
This test is very similar to the previous one. Can they be combined?
##########
clients/src/main/java/org/apache/kafka/clients/ClientUtils.java:
##########
@@ -238,7 +239,8 @@ public static NetworkClient
createNetworkClient(AbstractConfig config,
int
maxInFlightRequestsPerConnection,
Metadata metadata,
Sensor throttleTimeSensor,
- ClientTelemetrySender
clientTelemetrySender) {
+ ClientTelemetrySender
clientTelemetrySender,
+ Uuid clientInstanceId) {
Review Comment:
nit: I would put `clientInstanceId` higher up the list of arguments, just
before `metrics` and just after `clientId` (in the case where there is a client
ID). The client ID and client instance ID seems like a pair of related
identifiers to me.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]