Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
21 changes: 14 additions & 7 deletions src/rdkafka_mock_handlers.c
Original file line number Diff line number Diff line change
Expand Up @@ -2155,10 +2155,10 @@ rd_kafka_mock_handle_AddPartitionsToTxn(rd_kafka_mock_connection_t *mconn,
/* Epoch */
rd_kafka_buf_read_i16(rkbuf, &pid.epoch);
/* #Topics */
rd_kafka_buf_read_i32(rkbuf, &TopicsCnt);
rd_kafka_buf_read_arraycnt(rkbuf, &TopicsCnt, RD_KAFKAP_TOPICS_MAX);

/* Response: #Results */
rd_kafka_buf_write_i32(resp, TopicsCnt);
rd_kafka_buf_write_arraycnt(resp, TopicsCnt);

/* Inject error */
all_err = rd_kafka_mock_next_request_error(mconn, resp);
Expand All @@ -2183,9 +2183,10 @@ rd_kafka_mock_handle_AddPartitionsToTxn(rd_kafka_mock_connection_t *mconn,
rd_kafka_buf_write_kstr(resp, &Topic);

/* #Partitions */
rd_kafka_buf_read_i32(rkbuf, &PartsCnt);
rd_kafka_buf_read_arraycnt(rkbuf, &PartsCnt,
RD_KAFKAP_PARTITIONS_MAX);
/* Response: #Partitions */
rd_kafka_buf_write_i32(resp, PartsCnt);
rd_kafka_buf_write_arraycnt(resp, PartsCnt);

mtopic = rd_kafka_mock_topic_find_by_kstr(mcluster, &Topic);

Expand All @@ -2206,7 +2207,13 @@ rd_kafka_mock_handle_AddPartitionsToTxn(rd_kafka_mock_connection_t *mconn,

/* Response: ErrorCode */
rd_kafka_buf_write_i16(resp, err);

rd_kafka_buf_write_tags_empty(resp);
}

rd_kafka_buf_skip_tags(rkbuf);

Copy link
Preview

Copilot AI Aug 22, 2025

Choose a reason for hiding this comment

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

The tag skipping should be moved inside the partition loop (after line 2211) to properly handle tags for each partition entry, not after processing all partitions for a topic.

Suggested change
rd_kafka_buf_skip_tags(rkbuf);
}

Copilot uses AI. Check for mistakes.

rd_kafka_buf_write_tags_empty(resp);
}

rd_kafka_mock_connection_send_response(mconn, resp);
Expand Down Expand Up @@ -3013,12 +3020,12 @@ const struct rd_kafka_mock_api_handler
[RD_KAFKAP_LeaveGroup] = {0, 4, 4, rd_kafka_mock_handle_LeaveGroup},
[RD_KAFKAP_SyncGroup] = {0, 4, 4, rd_kafka_mock_handle_SyncGroup},
[RD_KAFKAP_AddPartitionsToTxn] =
{0, 1, -1, rd_kafka_mock_handle_AddPartitionsToTxn},
[RD_KAFKAP_AddOffsetsToTxn] = {0, 1, -1,
{0, 3, 3, rd_kafka_mock_handle_AddPartitionsToTxn},
[RD_KAFKAP_AddOffsetsToTxn] = {0, 3, 3,
rd_kafka_mock_handle_AddOffsetsToTxn},
[RD_KAFKAP_TxnOffsetCommit] = {0, 3, 3,
rd_kafka_mock_handle_TxnOffsetCommit},
[RD_KAFKAP_EndTxn] = {0, 1, -1, rd_kafka_mock_handle_EndTxn},
[RD_KAFKAP_EndTxn] = {0, 3, 3, rd_kafka_mock_handle_EndTxn},
[RD_KAFKAP_OffsetForLeaderEpoch] =
{2, 2, -1, rd_kafka_mock_handle_OffsetForLeaderEpoch},
[RD_KAFKAP_ConsumerGroupHeartbeat] =
Expand Down
35 changes: 20 additions & 15 deletions src/rdkafka_request.c
Original file line number Diff line number Diff line change
Expand Up @@ -6285,7 +6285,7 @@ rd_kafka_AddPartitionsToTxnRequest(rd_kafka_broker_t *rkb,
int TopicCnt = 0, PartCnt = 0;

ApiVersion = rd_kafka_broker_ApiVersion_supported(
rkb, RD_KAFKAP_AddPartitionsToTxn, 0, 0, NULL);
rkb, RD_KAFKAP_AddPartitionsToTxn, 0, 3, NULL);
if (ApiVersion == -1) {
rd_snprintf(errstr, errstr_size,
"AddPartitionsToTxnRequest (KIP-98) not supported "
Expand All @@ -6294,8 +6294,8 @@ rd_kafka_AddPartitionsToTxnRequest(rd_kafka_broker_t *rkb,
return RD_KAFKA_RESP_ERR__UNSUPPORTED_FEATURE;
}

rkbuf =
rd_kafka_buf_new_request(rkb, RD_KAFKAP_AddPartitionsToTxn, 1, 500);
rkbuf = rd_kafka_buf_new_flexver_request(
rkb, RD_KAFKAP_AddPartitionsToTxn, 1, 500, ApiVersion >= 3);

/* transactional_id */
rd_kafka_buf_write_str(rkbuf, transactional_id, -1);
Expand All @@ -6305,23 +6305,24 @@ rd_kafka_AddPartitionsToTxnRequest(rd_kafka_broker_t *rkb,
rd_kafka_buf_write_i16(rkbuf, pid.epoch);

/* Topics/partitions array (count updated later) */
of_TopicCnt = rd_kafka_buf_write_i32(rkbuf, 0);
of_TopicCnt = rd_kafka_buf_write_arraycnt_pos(rkbuf);

TAILQ_FOREACH(rktp, rktps, rktp_txnlink) {
if (last_rkt != rktp->rktp_rkt) {

if (last_rkt) {
/* Update last topic's partition count field */
rd_kafka_buf_update_i32(rkbuf, of_PartCnt,
PartCnt);
rd_kafka_buf_finalize_arraycnt(
rkbuf, of_PartCnt, PartCnt);
rd_kafka_buf_write_tags_empty(rkbuf);
of_PartCnt = -1;
}

/* Topic name */
rd_kafka_buf_write_kstr(rkbuf,
rktp->rktp_rkt->rkt_topic);
/* Partition count, updated later */
of_PartCnt = rd_kafka_buf_write_i32(rkbuf, 0);
of_PartCnt = rd_kafka_buf_write_arraycnt_pos(rkbuf);

PartCnt = 0;
TopicCnt++;
Expand All @@ -6334,9 +6335,12 @@ rd_kafka_AddPartitionsToTxnRequest(rd_kafka_broker_t *rkb,
}

/* Update last partition and topic count fields */
if (of_PartCnt != -1)
rd_kafka_buf_update_i32(rkbuf, (size_t)of_PartCnt, PartCnt);
rd_kafka_buf_update_i32(rkbuf, of_TopicCnt, TopicCnt);
if (of_PartCnt != -1) {
rd_kafka_buf_finalize_arraycnt(rkbuf, (size_t)of_PartCnt,
PartCnt);
rd_kafka_buf_write_tags_empty(rkbuf);
}
rd_kafka_buf_finalize_arraycnt(rkbuf, of_TopicCnt, TopicCnt);

rd_kafka_buf_ApiVersion_set(rkbuf, ApiVersion, 0);

Expand Down Expand Up @@ -6373,7 +6377,7 @@ rd_kafka_AddOffsetsToTxnRequest(rd_kafka_broker_t *rkb,
int16_t ApiVersion = 0;

ApiVersion = rd_kafka_broker_ApiVersion_supported(
rkb, RD_KAFKAP_AddOffsetsToTxn, 0, 0, NULL);
rkb, RD_KAFKAP_AddOffsetsToTxn, 0, 3, NULL);
if (ApiVersion == -1) {
rd_snprintf(errstr, errstr_size,
"AddOffsetsToTxnRequest (KIP-98) not supported "
Expand All @@ -6382,8 +6386,8 @@ rd_kafka_AddOffsetsToTxnRequest(rd_kafka_broker_t *rkb,
return RD_KAFKA_RESP_ERR__UNSUPPORTED_FEATURE;
}

rkbuf =
rd_kafka_buf_new_request(rkb, RD_KAFKAP_AddOffsetsToTxn, 1, 100);
rkbuf = rd_kafka_buf_new_flexver_request(rkb, RD_KAFKAP_AddOffsetsToTxn,
1, 100, ApiVersion >= 3);

/* transactional_id */
rd_kafka_buf_write_str(rkbuf, transactional_id, -1);
Expand Down Expand Up @@ -6428,7 +6432,7 @@ rd_kafka_resp_err_t rd_kafka_EndTxnRequest(rd_kafka_broker_t *rkb,
int16_t ApiVersion = 0;

ApiVersion = rd_kafka_broker_ApiVersion_supported(rkb, RD_KAFKAP_EndTxn,
0, 1, NULL);
0, 3, NULL);
if (ApiVersion == -1) {
rd_snprintf(errstr, errstr_size,
"EndTxnRequest (KIP-98) not supported "
Expand All @@ -6437,7 +6441,8 @@ rd_kafka_resp_err_t rd_kafka_EndTxnRequest(rd_kafka_broker_t *rkb,
return RD_KAFKA_RESP_ERR__UNSUPPORTED_FEATURE;
}

rkbuf = rd_kafka_buf_new_request(rkb, RD_KAFKAP_EndTxn, 1, 500);
rkbuf = rd_kafka_buf_new_flexver_request(rkb, RD_KAFKAP_EndTxn, 1, 500,
ApiVersion >= 3);

/* transactional_id */
rd_kafka_buf_write_str(rkbuf, transactional_id, -1);
Expand Down
8 changes: 6 additions & 2 deletions src/rdkafka_txnmgr.c
Original file line number Diff line number Diff line change
Expand Up @@ -842,7 +842,7 @@ static void rd_kafka_txn_handle_AddPartitionsToTxn(rd_kafka_t *rk,

rd_kafka_buf_read_throttle_time(rkbuf);

rd_kafka_buf_read_i32(rkbuf, &TopicCnt);
rd_kafka_buf_read_arraycnt(rkbuf, &TopicCnt, RD_KAFKAP_TOPICS_MAX);

while (TopicCnt-- > 0) {
rd_kafkap_str_t Topic;
Expand All @@ -851,7 +851,8 @@ static void rd_kafka_txn_handle_AddPartitionsToTxn(rd_kafka_t *rk,
rd_bool_t request_error = rd_false;

rd_kafka_buf_read_str(rkbuf, &Topic);
rd_kafka_buf_read_i32(rkbuf, &PartCnt);
rd_kafka_buf_read_arraycnt(rkbuf, &PartCnt,
RD_KAFKAP_PARTITIONS_MAX);

rkt = rd_kafka_topic_find0(rk, &Topic);
if (rkt)
Expand All @@ -865,6 +866,7 @@ static void rd_kafka_txn_handle_AddPartitionsToTxn(rd_kafka_t *rk,

rd_kafka_buf_read_i32(rkbuf, &Partition);
rd_kafka_buf_read_i16(rkbuf, &ErrorCode);
rd_kafka_buf_skip_tags(rkbuf);

if (rkt)
rktp = rd_kafka_toppar_get(rkt, Partition,
Expand Down Expand Up @@ -979,6 +981,8 @@ static void rd_kafka_txn_handle_AddPartitionsToTxn(rd_kafka_t *rk,
break; /* Request-level error seen, bail out */
}

rd_kafka_buf_skip_tags(rkbuf);

if (rkt) {
rd_kafka_topic_rdunlock(rkt);
rd_kafka_topic_destroy0(rkt);
Expand Down