Skip to content
Merged
Show file tree
Hide file tree
Changes from 5 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -201,7 +201,7 @@ public short oldestVersion() {

public List<Short> allVersions() {
List<Short> versions = new ArrayList<>(latestVersion() - oldestVersion() + 1);
for (short version = oldestVersion(); version < latestVersion(); version++) {
for (short version = oldestVersion(); version <= latestVersion(); version++) {

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is not a critical bug as the method allVersions is used by testing only.

versions.add(version);
}
return versions;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ public void testListOffsetsResponseVersions() throws Exception {
.setPartitions(Collections.singletonList(partition)));
Supplier<ListOffsetsResponseData> response = () -> new ListOffsetsResponseData()
.setTopics(topics);
for (short version = 0; version <= ApiKeys.LIST_OFFSETS.latestVersion(); version++) {
for (short version : ApiKeys.LIST_OFFSETS.allVersions()) {
ListOffsetsResponseData responseData = response.get();
if (version > 0) {
responseData.topics().get(0).partitions().get(0)
Expand Down Expand Up @@ -459,7 +459,7 @@ public void testOffsetCommitRequestVersions() throws Exception {
))))
.setRetentionTimeMs(20);

for (short version = 0; version <= ApiKeys.OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_COMMIT.allVersions()) {
OffsetCommitRequestData requestData = request.get();
if (version < 1) {
requestData.setMemberId("");
Expand All @@ -485,7 +485,7 @@ public void testOffsetCommitRequestVersions() throws Exception {
if (version == 1) {
testEquivalentMessageRoundTrip(version, requestData);
} else if (version >= 2 && version <= 4) {
testAllMessageRoundTripsBetweenVersions(version, (short) 4, requestData, requestData);
testAllMessageRoundTripsBetweenVersions(version, (short) 5, requestData, requestData);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why are we changing this?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

testAllMessageRoundTripsBetweenVersions exclude the end number so it should pass 5 so as to test version 4

} else {
testAllMessageRoundTripsFromVersion(version, requestData);
}
Expand All @@ -509,7 +509,7 @@ public void testOffsetCommitResponseVersions() throws Exception {
)
.setThrottleTimeMs(20);

for (short version = 0; version <= ApiKeys.OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_COMMIT.allVersions()) {
OffsetCommitResponseData responseData = response.get();
if (version < 3) {
responseData.setThrottleTimeMs(0);
Expand Down Expand Up @@ -568,7 +568,7 @@ public void testTxnOffsetCommitRequestVersions() throws Exception {
.setCommittedOffset(offset)
))));

for (short version = 0; version <= ApiKeys.TXN_OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.TXN_OFFSET_COMMIT.allVersions()) {
TxnOffsetCommitRequestData requestData = request.get();
if (version < 2) {
requestData.topics().get(0).partitions().get(0).setCommittedLeaderEpoch(-1);
Expand Down Expand Up @@ -632,7 +632,7 @@ public void testOffsetFetchVersions() throws Exception {
.setTopics(topics)
.setRequireStable(true);

for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
final short finalVersion = version;
if (version < 2) {
assertThrows(NullPointerException.class, () -> testAllMessageRoundTripsFromVersion(finalVersion, allPartitionData));
Expand Down Expand Up @@ -661,7 +661,7 @@ public void testOffsetFetchVersions() throws Exception {
.setErrorCode(Errors.UNKNOWN_TOPIC_OR_PARTITION.code())))))
.setErrorCode(Errors.NOT_COORDINATOR.code())
.setThrottleTimeMs(10);
for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
OffsetFetchResponseData responseData = response.get();
if (version <= 1) {
responseData.setErrorCode(Errors.NONE.code());
Expand Down Expand Up @@ -720,7 +720,7 @@ public void testProduceResponseVersions() throws Exception {
.setErrorMessage(errorMessage)))).iterator()))
.setThrottleTimeMs(throttleTimeMs);

for (short version = 0; version <= ApiKeys.PRODUCE.latestVersion(); version++) {
for (short version : ApiKeys.PRODUCE.allVersions()) {
ProduceResponseData responseData = response.get();

if (version < 8) {
Expand All @@ -741,9 +741,9 @@ public void testProduceResponseVersions() throws Exception {
}

if (version >= 3 && version <= 4) {
testAllMessageRoundTripsBetweenVersions(version, (short) 4, responseData, responseData);
testAllMessageRoundTripsBetweenVersions(version, (short) 5, responseData, responseData);
} else if (version >= 6 && version <= 7) {
testAllMessageRoundTripsBetweenVersions(version, (short) 7, responseData, responseData);
testAllMessageRoundTripsBetweenVersions(version, (short) 8, responseData, responseData);
} else {
testEquivalentMessageRoundTrip(version, responseData);
}
Expand Down Expand Up @@ -924,8 +924,8 @@ public void testNonIgnorableFieldWithDefaultNull() {
@Test
public void testWriteNullForNonNullableFieldRaisesException() {
CreateTopicsRequestData createTopics = new CreateTopicsRequestData().setTopics(null);
for (short i = (short) 0; i <= createTopics.highestSupportedVersion(); i++) {
verifyWriteRaisesNpe(i, createTopics);
for (short version : ApiKeys.CREATE_TOPICS.allVersions()) {
verifyWriteRaisesNpe(version, createTopics);
}
MetadataRequestData metadata = new MetadataRequestData().setTopics(null);
verifyWriteRaisesNpe((short) 0, metadata);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ public void testConstructor() {

AddPartitionsToTxnRequest.Builder builder = new AddPartitionsToTxnRequest.Builder(transactionalId, producerId, producerEpoch, partitions);

for (short version = 0; version <= ApiKeys.ADD_PARTITIONS_TO_TXN.latestVersion(); version++) {
for (short version : ApiKeys.ADD_PARTITIONS_TO_TXN.allVersions()) {
AddPartitionsToTxnRequest request = builder.build(version);

assertEquals(transactionalId, request.data().transactionalId());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,7 @@ public void testParse() {
.setThrottleTimeMs(throttleTimeMs);
AddPartitionsToTxnResponse response = new AddPartitionsToTxnResponse(data);

for (short version = 0; version <= ApiKeys.ADD_PARTITIONS_TO_TXN.latestVersion(); version++) {
for (short version : ApiKeys.ADD_PARTITIONS_TO_TXN.allVersions()) {
AddPartitionsToTxnResponse parsedResponse = AddPartitionsToTxnResponse.parse(response.serialize(version), version);
assertEquals(expectedErrorCounts, parsedResponse.errorCounts());
assertEquals(throttleTimeMs, parsedResponse.throttleTimeMs());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ public void testUnsupportedVersion() {

@Test
public void testGetErrorResponse() {
for (short version = CONTROLLED_SHUTDOWN.oldestVersion(); version < CONTROLLED_SHUTDOWN.latestVersion(); version++) {
for (short version : CONTROLLED_SHUTDOWN.allVersions()) {
ControlledShutdownRequest.Builder builder = new ControlledShutdownRequest.Builder(
new ControlledShutdownRequestData().setBrokerId(1), version);
ControlledShutdownRequest request = builder.build();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public void testConstructor() {
.setProducerId(producerId)
.setTransactionalId(transactionId));

for (short version = 0; version <= ApiKeys.END_TXN.latestVersion(); version++) {
for (short version : ApiKeys.END_TXN.allVersions()) {
EndTxnRequest request = builder.build(version);

EndTxnResponse response = request.getErrorResponse(throttleTimeMs, Errors.NOT_COORDINATOR.exception());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ public void testConstructor() {

Map<Errors, Integer> expectedErrorCounts = Collections.singletonMap(Errors.NOT_COORDINATOR, 1);

for (short version = 0; version <= ApiKeys.END_TXN.latestVersion(); version++) {
for (short version : ApiKeys.END_TXN.allVersions()) {
EndTxnResponse response = new EndTxnResponse(data);
assertEquals(expectedErrorCounts, response.errorCounts());
assertEquals(throttleTimeMs, response.throttleTimeMs());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ public void testGetPrincipal() {

@Test
public void testToSend() throws IOException {
for (short version = ApiKeys.ENVELOPE.oldestVersion(); version <= ApiKeys.ENVELOPE.latestVersion(); version++) {
for (short version : ApiKeys.ENVELOPE.allVersions()) {
ByteBuffer requestData = ByteBuffer.wrap("foobar".getBytes());
RequestHeader header = new RequestHeader(ApiKeys.ENVELOPE, version, "clientId", 15);
EnvelopeRequest request = new EnvelopeRequest.Builder(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ class EnvelopeResponseTest {

@Test
public void testToSend() {
for (short version = ApiKeys.ENVELOPE.oldestVersion(); version <= ApiKeys.ENVELOPE.latestVersion(); version++) {
for (short version : ApiKeys.ENVELOPE.allVersions()) {
ByteBuffer responseData = ByteBuffer.wrap("foobar".getBytes());
EnvelopeResponse response = new EnvelopeResponse(responseData, Errors.NONE);
short headerVersion = ApiKeys.ENVELOPE.responseHeaderVersion(version);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,7 @@ public void testGetErrorResponse() {
Uuid topicId = Uuid.randomUuid();
String topicName = "topic";
int partition = 0;

for (short version = LEADER_AND_ISR.oldestVersion(); version <= LEADER_AND_ISR.latestVersion(); version++) {
for (short version : LEADER_AND_ISR.allVersions()) {
LeaderAndIsrRequest request = new LeaderAndIsrRequest.Builder(version, 0, 0, 0,
Collections.singletonList(new LeaderAndIsrPartitionState()
.setTopicName(topicName)
Expand Down Expand Up @@ -108,7 +107,7 @@ public void testGetErrorResponse() {
*/
@Test
public void testVersionLogic() {
for (short version = LEADER_AND_ISR.oldestVersion(); version <= LEADER_AND_ISR.latestVersion(); version++) {
for (short version : LEADER_AND_ISR.allVersions()) {
List<LeaderAndIsrPartitionState> partitionStates = asList(
new LeaderAndIsrPartitionState()
.setTopicName("topic0")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ public void testErrorCountsFromGetErrorResponse() {

@Test
public void testErrorCountsWithTopLevelError() {
for (short version = LEADER_AND_ISR.oldestVersion(); version < LEADER_AND_ISR.latestVersion(); version++) {
for (short version : LEADER_AND_ISR.allVersions()) {
LeaderAndIsrResponse response;
if (version < 5) {
List<LeaderAndIsrPartitionError> partitions = createPartitions("foo",
Expand All @@ -92,7 +92,7 @@ public void testErrorCountsWithTopLevelError() {

@Test
public void testErrorCountsNoTopLevelError() {
for (short version = LEADER_AND_ISR.oldestVersion(); version < LEADER_AND_ISR.latestVersion(); version++) {
for (short version : LEADER_AND_ISR.allVersions()) {
LeaderAndIsrResponse response;
if (version < 5) {
List<LeaderAndIsrPartitionError> partitions = createPartitions("foo",
Expand All @@ -116,7 +116,7 @@ public void testErrorCountsNoTopLevelError() {

@Test
public void testToString() {
for (short version = LEADER_AND_ISR.oldestVersion(); version < LEADER_AND_ISR.latestVersion(); version++) {
for (short version : LEADER_AND_ISR.allVersions()) {
LeaderAndIsrResponse response;
if (version < 5) {
List<LeaderAndIsrPartitionError> partitions = createPartitions("foo",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ public void testMultiLeaveConstructor() {
.setGroupId(groupId)
.setMembers(members);

for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
try {
LeaveGroupRequest request = builder.build(version);
if (version <= 2) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ public void testConstructorWithMemberResponses() {
expectedErrorCounts.put(Errors.UNKNOWN_MEMBER_ID, 1);
expectedErrorCounts.put(Errors.FENCED_INSTANCE_ID, 1);

for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
LeaveGroupResponse leaveGroupResponse = new LeaveGroupResponse(memberResponses,
Errors.NONE,
throttleTimeMs,
Expand Down Expand Up @@ -95,7 +95,7 @@ public void testConstructorWithMemberResponses() {
@Test
public void testShouldThrottle() {
LeaveGroupResponse response = new LeaveGroupResponse(new LeaveGroupResponseData());
for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
if (version >= 2) {
assertTrue(response.shouldClientThrottle(version));
} else {
Expand All @@ -109,7 +109,7 @@ public void testEqualityWithSerialization() {
LeaveGroupResponseData responseData = new LeaveGroupResponseData()
.setErrorCode(Errors.NONE.code())
.setThrottleTimeMs(throttleTimeMs);
for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
LeaveGroupResponse primaryResponse = LeaveGroupResponse.parse(
MessageUtil.toByteBuffer(responseData, version), version);
LeaveGroupResponse secondaryResponse = LeaveGroupResponse.parse(
Expand All @@ -129,7 +129,7 @@ public void testParse() {
.setErrorCode(Errors.NOT_COORDINATOR.code())
.setThrottleTimeMs(throttleTimeMs);

for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
ByteBuffer buffer = MessageUtil.toByteBuffer(data, version);
LeaveGroupResponse leaveGroupResponse = LeaveGroupResponse.parse(buffer, version);
assertEquals(expectedErrorCounts, leaveGroupResponse.errorCounts());
Expand All @@ -146,7 +146,7 @@ public void testParse() {

@Test
public void testEqualityWithMemberResponses() {
for (short version = 0; version <= ApiKeys.LEAVE_GROUP.latestVersion(); version++) {
for (short version : ApiKeys.LEAVE_GROUP.allVersions()) {
List<MemberResponse> localResponses = version > 2 ? memberResponses : memberResponses.subList(0, 1);
LeaveGroupResponse primaryResponse = new LeaveGroupResponse(localResponses,
Errors.NONE,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ public void testConstructor() {

OffsetCommitRequest.Builder builder = new OffsetCommitRequest.Builder(data);

for (short version = 0; version <= ApiKeys.TXN_OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.TXN_OFFSET_COMMIT.allVersions()) {
OffsetCommitRequest request = builder.build(version);
assertEquals(expectedOffsets, request.offsets());

Expand Down Expand Up @@ -130,7 +130,7 @@ public void testVersionSupportForGroupInstanceId() {
.setGroupInstanceId(groupInstanceId)
);

for (short version = 0; version <= ApiKeys.OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_COMMIT.allVersions()) {
if (version >= 7) {
builder.build(version);
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ public void testParse() {
))
.setThrottleTimeMs(throttleTimeMs);

for (short version = 0; version <= ApiKeys.OFFSET_COMMIT.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_COMMIT.allVersions()) {
ByteBuffer buffer = MessageUtil.toByteBuffer(data, version);
OffsetCommitResponse response = OffsetCommitResponse.parse(buffer, version);
assertEquals(expectedErrorCounts, response.errorCounts());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ public void testConstructor() {
));
}

for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
OffsetFetchRequest request = builder.build(version);
assertFalse(request.isAllPartitions());
assertEquals(groupId, request.groupId());
Expand All @@ -101,7 +101,7 @@ public void testConstructor() {

@Test
public void testConstructorFailForUnsupportedRequireStable() {
for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
// The builder needs to be initialized every cycle as the internal data `requireStable` flag is flipped.
builder = new OffsetFetchRequest.Builder(groupId, true, null, false);
final short finalVersion = version;
Expand All @@ -123,7 +123,7 @@ public void testConstructorFailForUnsupportedRequireStable() {

@Test
public void testBuildThrowForUnsupportedRequireStable() {
for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
builder = new OffsetFetchRequest.Builder(groupId, true, null, true);
if (version < 7) {
final short finalVersion = version;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@ public void testStructBuild() {

OffsetFetchResponse latestResponse = new OffsetFetchResponse(throttleTimeMs, Errors.NONE, partitionDataMap);

for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
OffsetFetchResponseData data = new OffsetFetchResponseData(
new ByteBufferAccessor(latestResponse.serialize(version)), version);

Expand Down Expand Up @@ -154,7 +154,7 @@ public void testStructBuild() {
@Test
public void testShouldThrottle() {
OffsetFetchResponse response = new OffsetFetchResponse(throttleTimeMs, Errors.NONE, partitionDataMap);
for (short version = 0; version <= ApiKeys.OFFSET_FETCH.latestVersion(); version++) {
for (short version : ApiKeys.OFFSET_FETCH.allVersions()) {
if (version >= 4) {
assertTrue(response.shouldClientThrottle(version));
} else {
Expand Down
Loading