use std::collections::{BTreeMap, BTreeSet};
use crate::native::{
KafkaClientError, KafkaClientResult,
model::{PayloadBatchBuilder, TopicPartition},
protocol::{Decoder, Encoder, decode_record_batches},
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ApiKeyVersion {
pub(crate) api_key: i16,
pub(crate) min_version: i16,
pub(crate) max_version: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ApiVersionsResponse {
pub(crate) error_code: i16,
pub(crate) api_keys: Vec<ApiKeyVersion>,
}
pub(crate) fn encode_api_versions_request(
client_software_name: &str,
client_software_version: &str,
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(64);
encoder.put_compact_string(client_software_name)?;
encoder.put_compact_string(client_software_version)?;
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl ApiVersionsResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
let flexible = version >= 3;
let mut decoder = Decoder::new(bytes);
let error_code = decoder.get_i16()?;
let api_count = decoder.get_array_len(flexible)?;
let mut api_keys = Vec::with_capacity(api_count);
for _ in 0..api_count {
api_keys.push(ApiKeyVersion {
api_key: decoder.get_i16()?,
min_version: decoder.get_i16()?,
max_version: decoder.get_i16()?,
});
if flexible {
decoder.skip_tags()?;
}
}
if version >= 1 {
let _throttle_time_ms = decoder.get_i32()?;
}
if flexible {
decoder.skip_tags()?;
}
expect_done(&decoder, "ApiVersionsResponse")?;
Ok(Self {
error_code,
api_keys,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SaslHandshakeResponse {
pub(crate) error_code: i16,
pub(crate) mechanisms: Vec<String>,
}
pub(crate) fn encode_sasl_handshake_request(mechanism: &str) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(mechanism.len() + 2);
encoder.put_string(mechanism)?;
Ok(encoder.into_inner())
}
impl SaslHandshakeResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 1 {
return Err(KafkaClientError::unsupported(format!(
"SaslHandshakeResponse v{version}; native security requires v1"
)));
}
let mut decoder = Decoder::new(bytes);
let error_code = decoder.get_i16()?;
let count = decoder.get_array_len(false)?;
let mut mechanisms = Vec::with_capacity(count);
for _ in 0..count {
mechanisms.push(decoder.get_string()?);
}
expect_done(&decoder, "SaslHandshakeResponse")?;
Ok(Self {
error_code,
mechanisms,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SaslAuthenticateResponse {
pub(crate) error_code: i16,
pub(crate) error_message: Option<String>,
pub(crate) auth_bytes: Vec<u8>,
pub(crate) session_lifetime_ms: i64,
}
pub(crate) fn encode_sasl_authenticate_request(auth_bytes: &[u8]) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(auth_bytes.len() + 4);
encoder.put_bytes(auth_bytes)?;
Ok(encoder.into_inner())
}
impl SaslAuthenticateResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if !matches!(version, 0 | 1) {
return Err(KafkaClientError::unsupported(format!(
"SaslAuthenticateResponse v{version}; native security supports v0/v1"
)));
}
let mut decoder = Decoder::new(bytes);
let error_code = decoder.get_i16()?;
let error_message = decoder.get_nullable_string()?;
let auth_bytes = decoder
.get_nullable_bytes()?
.ok_or_else(|| KafkaClientError::protocol("SaslAuthenticate auth_bytes was null"))?
.to_vec();
let session_lifetime_ms = if version >= 1 { decoder.get_i64()? } else { 0 };
expect_done(&decoder, "SaslAuthenticateResponse")?;
Ok(Self {
error_code,
error_message,
auth_bytes,
session_lifetime_ms,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct MetadataResponse {
pub(crate) brokers: Vec<MetadataBroker>,
pub(crate) cluster_id: Option<String>,
pub(crate) topics: Vec<MetadataTopic>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct MetadataBroker {
pub(crate) node_id: i32,
pub(crate) host: String,
pub(crate) port: i32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct MetadataTopic {
pub(crate) error_code: i16,
pub(crate) name: String,
pub(crate) partitions: Vec<MetadataPartition>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct MetadataPartition {
pub(crate) error_code: i16,
pub(crate) partition_index: i32,
pub(crate) leader_id: i32,
pub(crate) leader_epoch: i32,
}
pub(crate) fn encode_metadata_request(topics: &[String]) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(128);
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(topic)?;
encoder.put_empty_tags();
}
encoder.put_bool(false);
encoder.put_bool(false);
encoder.put_bool(false);
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl MetadataResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 9 {
return Err(KafkaClientError::unsupported(format!(
"MetadataResponse v{version}; KC-1 implements v9"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let broker_count = decoder.get_array_len(true)?;
let mut brokers = Vec::with_capacity(broker_count);
for _ in 0..broker_count {
brokers.push(MetadataBroker {
node_id: decoder.get_i32()?,
host: decoder.get_compact_string()?,
port: decoder.get_i32()?,
});
let _rack = decoder.get_compact_nullable_string()?;
decoder.skip_tags()?;
}
let cluster_id = decoder.get_compact_nullable_string()?;
let _controller_id = decoder.get_i32()?;
let topic_count = decoder.get_array_len(true)?;
let mut topics = Vec::with_capacity(topic_count);
for _ in 0..topic_count {
let error_code = decoder.get_i16()?;
let name = decoder.get_compact_string()?;
let _is_internal = decoder.get_bool()?;
let partition_count = decoder.get_array_len(true)?;
let mut partitions = Vec::with_capacity(partition_count);
for _ in 0..partition_count {
let error_code = decoder.get_i16()?;
let partition_index = decoder.get_i32()?;
let leader_id = decoder.get_i32()?;
let leader_epoch = decoder.get_i32()?;
skip_i32_array(&mut decoder, true)?;
skip_i32_array(&mut decoder, true)?;
skip_i32_array(&mut decoder, true)?;
decoder.skip_tags()?;
partitions.push(MetadataPartition {
error_code,
partition_index,
leader_id,
leader_epoch,
});
}
let _topic_authorized_operations = decoder.get_i32()?;
decoder.skip_tags()?;
topics.push(MetadataTopic {
error_code,
name,
partitions,
});
}
let _cluster_authorized_operations = decoder.get_i32()?;
decoder.skip_tags()?;
expect_done(&decoder, "MetadataResponse")?;
Ok(Self {
brokers,
cluster_id,
topics,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ListOffsetTopicRequest {
pub(crate) name: String,
pub(crate) partitions: Vec<ListOffsetPartitionRequest>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ListOffsetPartitionRequest {
pub(crate) partition_index: i32,
pub(crate) current_leader_epoch: i32,
pub(crate) timestamp: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ListOffsetResult {
pub(crate) topic: String,
pub(crate) partition_index: i32,
pub(crate) error_code: i16,
pub(crate) offset: i64,
}
pub(crate) fn encode_list_offsets_request(
topics: &[ListOffsetTopicRequest],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(128);
encoder.put_i32(-1);
encoder.put_i8(0);
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(&topic.name)?;
encoder.put_array_len(topic.partitions.len(), true)?;
for partition in &topic.partitions {
encoder.put_i32(partition.partition_index);
encoder.put_i32(partition.current_leader_epoch);
encoder.put_i64(partition.timestamp);
encoder.put_empty_tags();
}
encoder.put_empty_tags();
}
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
pub(crate) fn decode_list_offsets_response(
version: i16,
bytes: &[u8],
) -> KafkaClientResult<Vec<ListOffsetResult>> {
if version != 6 {
return Err(KafkaClientError::unsupported(format!(
"ListOffsetsResponse v{version}; KC-1 implements v6"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let topic_count = decoder.get_array_len(true)?;
let mut results = Vec::new();
for _ in 0..topic_count {
let topic = decoder.get_compact_string()?;
let partition_count = decoder.get_array_len(true)?;
results.reserve(partition_count);
for _ in 0..partition_count {
results.push(ListOffsetResult {
topic: topic.clone(),
partition_index: decoder.get_i32()?,
error_code: decoder.get_i16()?,
offset: {
let _timestamp = decoder.get_i64()?;
decoder.get_i64()?
},
});
let _leader_epoch = decoder.get_i32()?;
decoder.skip_tags()?;
}
decoder.skip_tags()?;
}
decoder.skip_tags()?;
expect_done(&decoder, "ListOffsetsResponse")?;
Ok(results)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FindCoordinatorResponse {
pub(crate) error_code: i16,
pub(crate) node_id: i32,
pub(crate) host: String,
pub(crate) port: i32,
}
pub(crate) fn encode_find_coordinator_request(group_id: &str) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(64);
encoder.put_compact_string(group_id)?;
encoder.put_i8(0);
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl FindCoordinatorResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 3 {
return Err(KafkaClientError::unsupported(format!(
"FindCoordinatorResponse v{version}; KC-2 implements v3"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let _error_message = decoder.get_compact_nullable_string()?;
let node_id = decoder.get_i32()?;
let host = decoder.get_compact_string()?;
let port = decoder.get_i32()?;
decoder.skip_tags()?;
expect_done(&decoder, "FindCoordinatorResponse")?;
Ok(Self {
error_code,
node_id,
host,
port,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetCommitTopicRequest {
pub(crate) name: String,
pub(crate) partitions: Vec<OffsetCommitPartitionRequest>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetCommitPartitionRequest {
pub(crate) partition_index: i32,
pub(crate) committed_offset: i64,
pub(crate) committed_leader_epoch: i32,
pub(crate) committed_metadata: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetCommitResult {
pub(crate) topic: String,
pub(crate) partition_index: i32,
pub(crate) error_code: i16,
}
pub(crate) fn encode_offset_commit_request(
group_id: &str,
generation_id: i32,
member_id: &str,
group_instance_id: Option<&str>,
topics: &[OffsetCommitTopicRequest],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(256);
encoder.put_compact_string(group_id)?;
encoder.put_i32(generation_id);
encoder.put_compact_string(member_id)?;
encoder.put_compact_nullable_string(group_instance_id)?;
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(&topic.name)?;
encoder.put_array_len(topic.partitions.len(), true)?;
for partition in &topic.partitions {
encoder.put_i32(partition.partition_index);
encoder.put_i64(partition.committed_offset);
encoder.put_i32(partition.committed_leader_epoch);
encoder.put_compact_nullable_string(partition.committed_metadata.as_deref())?;
encoder.put_empty_tags();
}
encoder.put_empty_tags();
}
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
pub(crate) fn decode_offset_commit_response(
version: i16,
bytes: &[u8],
) -> KafkaClientResult<Vec<OffsetCommitResult>> {
if version != 8 {
return Err(KafkaClientError::unsupported(format!(
"OffsetCommitResponse v{version}; KC-2 implements v8"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let topic_count = decoder.get_array_len(true)?;
let mut results = Vec::new();
for _ in 0..topic_count {
let topic = decoder.get_compact_string()?;
let partition_count = decoder.get_array_len(true)?;
results.reserve(partition_count);
for _ in 0..partition_count {
results.push(OffsetCommitResult {
topic: topic.clone(),
partition_index: decoder.get_i32()?,
error_code: decoder.get_i16()?,
});
decoder.skip_tags()?;
}
decoder.skip_tags()?;
}
decoder.skip_tags()?;
expect_done(&decoder, "OffsetCommitResponse")?;
Ok(results)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetFetchTopicRequest {
pub(crate) name: String,
pub(crate) partitions: Vec<i32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetFetchResponse {
pub(crate) error_code: i16,
pub(crate) offsets: Vec<OffsetFetchResult>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OffsetFetchResult {
pub(crate) topic: String,
pub(crate) partition_index: i32,
pub(crate) committed_offset: i64,
pub(crate) error_code: i16,
}
pub(crate) fn encode_offset_fetch_request(
group_id: &str,
topics: &[OffsetFetchTopicRequest],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(128);
encoder.put_compact_string(group_id)?;
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(&topic.name)?;
encoder.put_array_len(topic.partitions.len(), true)?;
for partition in &topic.partitions {
encoder.put_i32(*partition);
}
encoder.put_empty_tags();
}
encoder.put_bool(false);
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl OffsetFetchResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 7 {
return Err(KafkaClientError::unsupported(format!(
"OffsetFetchResponse v{version}; KC-2 implements v7"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let topic_count = decoder.get_array_len(true)?;
let mut offsets = Vec::new();
for _ in 0..topic_count {
let topic = decoder.get_compact_string()?;
let partition_count = decoder.get_array_len(true)?;
offsets.reserve(partition_count);
for _ in 0..partition_count {
let partition_index = decoder.get_i32()?;
let committed_offset = decoder.get_i64()?;
let _committed_leader_epoch = decoder.get_i32()?;
let _metadata = decoder.get_compact_nullable_string()?;
let error_code = decoder.get_i16()?;
decoder.skip_tags()?;
offsets.push(OffsetFetchResult {
topic: topic.clone(),
partition_index,
committed_offset,
error_code,
});
}
decoder.skip_tags()?;
}
let error_code = decoder.get_i16()?;
decoder.skip_tags()?;
expect_done(&decoder, "OffsetFetchResponse")?;
Ok(Self {
error_code,
offsets,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ConsumerGroupProtocol {
pub(crate) name: String,
pub(crate) metadata: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct JoinGroupResponse {
pub(crate) error_code: i16,
pub(crate) generation_id: i32,
pub(crate) protocol_name: String,
pub(crate) leader: String,
pub(crate) member_id: String,
pub(crate) members: Vec<JoinGroupMember>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct JoinGroupMember {
pub(crate) member_id: String,
pub(crate) group_instance_id: Option<String>,
pub(crate) metadata: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SyncGroupAssignment {
pub(crate) member_id: String,
pub(crate) assignment: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SyncGroupResponse {
pub(crate) error_code: i16,
pub(crate) assignment: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct HeartbeatResponse {
pub(crate) error_code: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct LeaveGroupResponse {
pub(crate) error_code: i16,
pub(crate) members: Vec<LeaveGroupMemberResponse>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct LeaveGroupMemberResponse {
pub(crate) member_id: String,
pub(crate) group_instance_id: Option<String>,
pub(crate) error_code: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ConsumerProtocolSubscription {
pub(crate) topics: Vec<String>,
pub(crate) owned_partitions: Vec<TopicPartition>,
pub(crate) generation_id: i32,
}
pub(crate) fn encode_join_group_request(
group_id: &str,
session_timeout_ms: i32,
rebalance_timeout_ms: i32,
member_id: &str,
group_instance_id: Option<&str>,
protocols: &[ConsumerGroupProtocol],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(256);
encoder.put_compact_string(group_id)?;
encoder.put_i32(session_timeout_ms);
encoder.put_i32(rebalance_timeout_ms);
encoder.put_compact_string(member_id)?;
encoder.put_compact_nullable_string(group_instance_id)?;
encoder.put_compact_string("consumer")?;
encoder.put_array_len(protocols.len(), true)?;
for protocol in protocols {
encoder.put_compact_string(&protocol.name)?;
encoder.put_compact_bytes(&protocol.metadata)?;
encoder.put_empty_tags();
}
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl JoinGroupResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 6 {
return Err(KafkaClientError::unsupported(format!(
"JoinGroupResponse v{version}; KC-3 implements v6"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let generation_id = decoder.get_i32()?;
let protocol_name = decoder.get_compact_string()?;
let leader = decoder.get_compact_string()?;
let member_id = decoder.get_compact_string()?;
let member_count = decoder.get_array_len(true)?;
let mut members = Vec::with_capacity(member_count);
for _ in 0..member_count {
members.push(JoinGroupMember {
member_id: decoder.get_compact_string()?,
group_instance_id: decoder.get_compact_nullable_string()?,
metadata: decoder.get_compact_bytes()?.to_vec(),
});
decoder.skip_tags()?;
}
decoder.skip_tags()?;
expect_done(&decoder, "JoinGroupResponse")?;
Ok(Self {
error_code,
generation_id,
protocol_name,
leader,
member_id,
members,
})
}
}
pub(crate) fn encode_sync_group_request(
group_id: &str,
generation_id: i32,
member_id: &str,
group_instance_id: Option<&str>,
assignments: &[SyncGroupAssignment],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(256);
encoder.put_compact_string(group_id)?;
encoder.put_i32(generation_id);
encoder.put_compact_string(member_id)?;
encoder.put_compact_nullable_string(group_instance_id)?;
encoder.put_array_len(assignments.len(), true)?;
for assignment in assignments {
encoder.put_compact_string(&assignment.member_id)?;
encoder.put_compact_bytes(&assignment.assignment)?;
encoder.put_empty_tags();
}
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl SyncGroupResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 4 {
return Err(KafkaClientError::unsupported(format!(
"SyncGroupResponse v{version}; KC-3 implements v4"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let assignment = decoder.get_compact_bytes()?.to_vec();
decoder.skip_tags()?;
expect_done(&decoder, "SyncGroupResponse")?;
Ok(Self {
error_code,
assignment,
})
}
}
pub(crate) fn encode_heartbeat_request(
group_id: &str,
generation_id: i32,
member_id: &str,
group_instance_id: Option<&str>,
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(96);
encoder.put_compact_string(group_id)?;
encoder.put_i32(generation_id);
encoder.put_compact_string(member_id)?;
encoder.put_compact_nullable_string(group_instance_id)?;
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl HeartbeatResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 4 {
return Err(KafkaClientError::unsupported(format!(
"HeartbeatResponse v{version}; KC-3 implements v4"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
decoder.skip_tags()?;
expect_done(&decoder, "HeartbeatResponse")?;
Ok(Self { error_code })
}
}
pub(crate) fn encode_leave_group_request(
group_id: &str,
member_id: &str,
group_instance_id: Option<&str>,
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(96);
encoder.put_compact_string(group_id)?;
encoder.put_array_len(1, true)?;
encoder.put_compact_string(member_id)?;
encoder.put_compact_nullable_string(group_instance_id)?;
encoder.put_empty_tags();
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl LeaveGroupResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 4 {
return Err(KafkaClientError::unsupported(format!(
"LeaveGroupResponse v{version}; KC-3 implements v4"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let member_count = decoder.get_array_len(true)?;
let mut members = Vec::with_capacity(member_count);
for _ in 0..member_count {
members.push(LeaveGroupMemberResponse {
member_id: decoder.get_compact_string()?,
group_instance_id: decoder.get_compact_nullable_string()?,
error_code: decoder.get_i16()?,
});
decoder.skip_tags()?;
}
decoder.skip_tags()?;
expect_done(&decoder, "LeaveGroupResponse")?;
Ok(Self {
error_code,
members,
})
}
}
pub(crate) fn encode_consumer_protocol_subscription(
topics: &[String],
owned_partitions: &[TopicPartition],
generation_id: i32,
) -> KafkaClientResult<Vec<u8>> {
let mut topics = topics.to_vec();
topics.sort();
topics.dedup();
let mut user_data = Encoder::with_capacity(4);
user_data.put_i32(generation_id);
let user_data = user_data.into_inner();
let mut encoder = Encoder::with_capacity(128);
encoder.put_i16(3);
encoder.put_array_len(topics.len(), false)?;
for topic in &topics {
encoder.put_string(topic)?;
}
encoder.put_nullable_bytes(Some(&user_data))?;
put_topic_partitions(&mut encoder, owned_partitions)?;
encoder.put_i32(generation_id);
encoder.put_nullable_string(None)?;
Ok(encoder.into_inner())
}
pub(crate) fn decode_consumer_protocol_subscription(
bytes: &[u8],
) -> KafkaClientResult<ConsumerProtocolSubscription> {
let mut decoder = Decoder::new(bytes);
let version = decoder.get_i16()?;
if !(0..=3).contains(&version) {
return Err(KafkaClientError::unsupported(format!(
"ConsumerProtocolSubscription v{version}; KC-3 implements v0..v3"
)));
}
let topic_count = decoder.get_array_len(false)?;
let mut topics = Vec::with_capacity(topic_count);
for _ in 0..topic_count {
topics.push(decoder.get_string()?);
}
let _user_data = decoder.get_nullable_bytes()?;
let owned_partitions = if version >= 1 {
get_topic_partitions(&mut decoder)?
} else {
Vec::new()
};
let generation_id = if version >= 2 { decoder.get_i32()? } else { -1 };
if version >= 3 {
let _rack_id = decoder.get_nullable_string()?;
}
expect_done(&decoder, "ConsumerProtocolSubscription")?;
Ok(ConsumerProtocolSubscription {
topics,
owned_partitions,
generation_id,
})
}
pub(crate) fn encode_consumer_protocol_assignment(
partitions: &[TopicPartition],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(128);
encoder.put_i16(3);
put_topic_partitions(&mut encoder, partitions)?;
encoder.put_nullable_bytes(None)?;
Ok(encoder.into_inner())
}
pub(crate) fn decode_consumer_protocol_assignment(
bytes: &[u8],
) -> KafkaClientResult<Vec<TopicPartition>> {
let mut decoder = Decoder::new(bytes);
let version = decoder.get_i16()?;
if !(0..=3).contains(&version) {
return Err(KafkaClientError::unsupported(format!(
"ConsumerProtocolAssignment v{version}; KC-3 implements v0..v3"
)));
}
let partitions = get_topic_partitions(&mut decoder)?;
let _user_data = decoder.get_nullable_bytes()?;
expect_done(&decoder, "ConsumerProtocolAssignment")?;
Ok(partitions)
}
fn put_topic_partitions(
encoder: &mut Encoder,
partitions: &[TopicPartition],
) -> KafkaClientResult<()> {
let mut grouped = BTreeMap::<String, BTreeSet<i32>>::new();
for partition in partitions {
grouped
.entry(partition.topic.clone())
.or_default()
.insert(partition.partition);
}
encoder.put_array_len(grouped.len(), false)?;
for (topic, partitions) in grouped {
encoder.put_string(&topic)?;
encoder.put_array_len(partitions.len(), false)?;
for partition in partitions {
encoder.put_i32(partition);
}
}
Ok(())
}
fn get_topic_partitions(decoder: &mut Decoder<'_>) -> KafkaClientResult<Vec<TopicPartition>> {
let topic_count = decoder.get_array_len(false)?;
let mut partitions = Vec::new();
for _ in 0..topic_count {
let topic = decoder.get_string()?;
let partition_count = decoder.get_array_len(false)?;
partitions.reserve(partition_count);
for _ in 0..partition_count {
partitions.push(TopicPartition::new(topic.clone(), decoder.get_i32()?));
}
}
Ok(partitions)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct InitProducerIdResponse {
pub(crate) error_code: i16,
pub(crate) producer_id: i64,
pub(crate) producer_epoch: i16,
}
pub(crate) fn encode_init_producer_id_request(
transactional_id: Option<&str>,
transaction_timeout_ms: i32,
producer_id: i64,
producer_epoch: i16,
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(32);
encoder.put_compact_nullable_string(transactional_id)?;
encoder.put_i32(transaction_timeout_ms);
encoder.put_i64(producer_id);
encoder.put_i16(producer_epoch);
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl InitProducerIdResponse {
pub(crate) fn decode(version: i16, bytes: &[u8]) -> KafkaClientResult<Self> {
if version != 4 {
return Err(KafkaClientError::unsupported(format!(
"InitProducerIdResponse v{version}; KC-11 implements v4"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let producer_id = decoder.get_i64()?;
let producer_epoch = decoder.get_i16()?;
decoder.skip_tags()?;
expect_done(&decoder, "InitProducerIdResponse")?;
Ok(Self {
error_code,
producer_id,
producer_epoch,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ProduceTopicRequest {
pub(crate) name: String,
pub(crate) partitions: Vec<ProducePartitionRequest>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ProducePartitionRequest {
pub(crate) index: i32,
pub(crate) records: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ProduceResult {
pub(crate) topic: String,
pub(crate) partition: i32,
pub(crate) error_code: i16,
pub(crate) base_offset: i64,
pub(crate) error_message: Option<String>,
}
pub(crate) fn encode_produce_request(
transactional_id: Option<&str>,
acks: i16,
timeout_ms: i32,
topics: &[ProduceTopicRequest],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(1024);
encoder.put_compact_nullable_string(transactional_id)?;
encoder.put_i16(acks);
encoder.put_i32(timeout_ms);
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(&topic.name)?;
encoder.put_array_len(topic.partitions.len(), true)?;
for partition in &topic.partitions {
encoder.put_i32(partition.index);
encoder.put_compact_bytes(&partition.records)?;
encoder.put_empty_tags();
}
encoder.put_empty_tags();
}
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
pub(crate) fn decode_produce_response(
version: i16,
bytes: &[u8],
) -> KafkaClientResult<Vec<ProduceResult>> {
if version != 9 {
return Err(KafkaClientError::unsupported(format!(
"ProduceResponse v{version}; KC-5 implements v9"
)));
}
let mut decoder = Decoder::new(bytes);
let topic_count = decoder.get_array_len(true)?;
let mut results = Vec::new();
for _ in 0..topic_count {
let topic = decoder.get_compact_string()?;
let partition_count = decoder.get_array_len(true)?;
results.reserve(partition_count);
for _ in 0..partition_count {
let partition = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let base_offset = decoder.get_i64()?;
let _log_append_time_ms = decoder.get_i64()?;
let _log_start_offset = decoder.get_i64()?;
let record_error_count = decoder.get_array_len(true)?;
for _ in 0..record_error_count {
let _batch_index = decoder.get_i32()?;
let _batch_error_message = decoder.get_compact_nullable_string()?;
decoder.skip_tags()?;
}
let error_message = decoder.get_compact_nullable_string()?;
decoder.skip_tags()?;
results.push(ProduceResult {
topic: topic.clone(),
partition,
error_code,
base_offset,
error_message,
});
}
decoder.skip_tags()?;
}
let _throttle_time_ms = decoder.get_i32()?;
decoder.skip_tags()?;
expect_done(&decoder, "ProduceResponse")?;
Ok(results)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchTopicRequest {
pub(crate) topic: String,
pub(crate) partitions: Vec<FetchPartitionRequest>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchPartitionRequest {
pub(crate) partition: i32,
pub(crate) current_leader_epoch: i32,
pub(crate) fetch_offset: i64,
pub(crate) last_fetched_epoch: i32,
pub(crate) log_start_offset: i64,
pub(crate) partition_max_bytes: i32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchResponse<'a> {
pub(crate) error_code: i16,
pub(crate) topics: Vec<FetchTopicResponse<'a>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchTopicResponse<'a> {
pub(crate) topic: String,
pub(crate) partitions: Vec<FetchPartitionResponse<'a>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchPartitionResponse<'a> {
pub(crate) partition_index: i32,
pub(crate) error_code: i16,
pub(crate) high_watermark: i64,
pub(crate) log_start_offset: i64,
pub(crate) records: Option<&'a [u8]>,
}
pub(crate) fn encode_fetch_request(
max_wait_ms: i32,
min_bytes: i32,
max_bytes: i32,
topics: &[FetchTopicRequest],
) -> KafkaClientResult<Vec<u8>> {
let mut encoder = Encoder::with_capacity(256);
encoder.put_i32(-1);
encoder.put_i32(max_wait_ms);
encoder.put_i32(min_bytes);
encoder.put_i32(max_bytes);
encoder.put_i8(0);
encoder.put_i32(0);
encoder.put_i32(-1);
encoder.put_array_len(topics.len(), true)?;
for topic in topics {
encoder.put_compact_string(&topic.topic)?;
encoder.put_array_len(topic.partitions.len(), true)?;
for partition in &topic.partitions {
encoder.put_i32(partition.partition);
encoder.put_i32(partition.current_leader_epoch);
encoder.put_i64(partition.fetch_offset);
encoder.put_i32(partition.last_fetched_epoch);
encoder.put_i64(partition.log_start_offset);
encoder.put_i32(partition.partition_max_bytes);
encoder.put_empty_tags();
}
encoder.put_empty_tags();
}
encoder.put_array_len(0, true)?;
encoder.put_compact_string("")?;
encoder.put_empty_tags();
Ok(encoder.into_inner())
}
impl<'a> FetchResponse<'a> {
pub(crate) fn decode(version: i16, bytes: &'a [u8]) -> KafkaClientResult<Self> {
if version != 12 {
return Err(KafkaClientError::unsupported(format!(
"FetchResponse v{version}; KC-1 implements v12"
)));
}
let mut decoder = Decoder::new(bytes);
let _throttle_time_ms = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let _session_id = decoder.get_i32()?;
let topic_count = decoder.get_array_len(true)?;
let mut topics = Vec::with_capacity(topic_count);
for _ in 0..topic_count {
let topic = decoder.get_compact_string()?;
let partition_count = decoder.get_array_len(true)?;
let mut partitions = Vec::with_capacity(partition_count);
for _ in 0..partition_count {
let partition_index = decoder.get_i32()?;
let error_code = decoder.get_i16()?;
let high_watermark = decoder.get_i64()?;
let _last_stable_offset = decoder.get_i64()?;
let log_start_offset = decoder.get_i64()?;
skip_aborted_transactions(&mut decoder)?;
let _preferred_read_replica = decoder.get_i32()?;
let records = decoder.get_compact_nullable_bytes()?;
decoder.skip_tags()?;
partitions.push(FetchPartitionResponse {
partition_index,
error_code,
high_watermark,
log_start_offset,
records,
});
}
decoder.skip_tags()?;
topics.push(FetchTopicResponse { topic, partitions });
}
decoder.skip_tags()?;
expect_done(&decoder, "FetchResponse")?;
Ok(Self { error_code, topics })
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FetchPartitionOutcome {
pub(crate) topic: String,
pub(crate) partition_index: i32,
pub(crate) error_code: i16,
pub(crate) high_watermark: i64,
pub(crate) log_start_offset: i64,
pub(crate) max_offset: Option<i64>,
}
#[derive(Debug)]
pub(crate) struct FetchBodyDecoder {
offsets: BTreeMap<(String, i32), i64>,
builder: PayloadBatchBuilder,
state: FetchDecodeState,
top_error_code: i16,
outcomes: Vec<FetchPartitionOutcome>,
}
#[derive(Debug)]
enum FetchDecodeState {
Header,
Topics {
topics_left: usize,
},
Partitions {
topics_left: usize,
topic: String,
parts_left: usize,
},
Blob {
topics_left: usize,
topic: String,
parts_left: usize,
partition: i32,
blob_len: usize,
decode: Option<i64>,
outcome_index: usize,
},
PartitionTags {
topics_left: usize,
topic: String,
parts_left: usize,
},
TopicTags {
topics_left: usize,
},
Trailer,
Done,
}
impl FetchBodyDecoder {
pub(crate) fn new(
version: i16,
offsets: BTreeMap<(String, i32), i64>,
builder: PayloadBatchBuilder,
) -> KafkaClientResult<Self> {
if version != 12 {
return Err(KafkaClientError::unsupported(format!(
"FetchResponse v{version}; KC-1 implements v12"
)));
}
Ok(Self {
offsets,
builder,
state: FetchDecodeState::Header,
top_error_code: 0,
outcomes: Vec::new(),
})
}
pub(crate) fn feed(&mut self, buf: &[u8], more: bool) -> KafkaClientResult<usize> {
let consumed = self.feed_inner(buf)?;
if !more && !matches!(self.state, FetchDecodeState::Done) {
return Err(KafkaClientError::protocol(
"truncated Kafka Fetch response body",
));
}
Ok(consumed)
}
pub(crate) fn finish(&self) -> KafkaClientResult<()> {
if matches!(self.state, FetchDecodeState::Done) {
Ok(())
} else {
Err(KafkaClientError::protocol(
"incomplete Kafka Fetch response body",
))
}
}
pub(crate) fn into_parts(self) -> (PayloadBatchBuilder, i16, Vec<FetchPartitionOutcome>) {
(self.builder, self.top_error_code, self.outcomes)
}
fn feed_inner(&mut self, buf: &[u8]) -> KafkaClientResult<usize> {
let mut pos = 0usize;
loop {
let state = std::mem::replace(&mut self.state, FetchDecodeState::Done);
let mut reader = ChunkReader::new(buf, pos);
match state {
FetchDecodeState::Header => {
let (Some(_throttle), Some(error_code), Some(_session)) =
(reader.get_i32(), reader.get_i16(), reader.get_i32())
else {
self.state = FetchDecodeState::Header;
return Ok(pos);
};
let Some(topic_count) = reader.get_compact_array_len()? else {
self.state = FetchDecodeState::Header;
return Ok(pos);
};
self.top_error_code = error_code;
self.state = FetchDecodeState::Topics {
topics_left: topic_count,
};
pos = reader.position();
}
FetchDecodeState::Topics { topics_left } => {
if topics_left == 0 {
self.state = FetchDecodeState::Trailer;
continue;
}
let Some(topic) = reader.get_compact_string()? else {
self.state = FetchDecodeState::Topics { topics_left };
return Ok(pos);
};
let Some(parts_left) = reader.get_compact_array_len()? else {
self.state = FetchDecodeState::Topics { topics_left };
return Ok(pos);
};
self.state = FetchDecodeState::Partitions {
topics_left,
topic,
parts_left,
};
pos = reader.position();
}
FetchDecodeState::Partitions {
topics_left,
topic,
parts_left,
} => {
if parts_left == 0 {
self.state = FetchDecodeState::TopicTags { topics_left };
continue;
}
macro_rules! need {
() => {{
self.state = FetchDecodeState::Partitions {
topics_left,
topic,
parts_left,
};
return Ok(pos);
}};
}
let (Some(partition_index), Some(error_code), Some(high_watermark)) =
(reader.get_i32(), reader.get_i16(), reader.get_i64())
else {
need!();
};
let (Some(_last_stable), Some(log_start_offset)) =
(reader.get_i64(), reader.get_i64())
else {
need!();
};
match reader.skip_aborted_transactions()? {
ChunkStep::Ready(()) => {}
ChunkStep::NeedMore => need!(),
}
let Some(_preferred_replica) = reader.get_i32() else {
need!();
};
let records_len = match reader.get_compact_nullable_bytes_len()? {
ChunkStep::Ready(len) => len,
ChunkStep::NeedMore => need!(),
};
let outcome_index = self.outcomes.len();
self.outcomes.push(FetchPartitionOutcome {
topic: topic.clone(),
partition_index,
error_code,
high_watermark,
log_start_offset,
max_offset: None,
});
let decode = if error_code == 0 {
self.offsets.get(&(topic.clone(), partition_index)).copied()
} else {
None
};
pos = reader.position();
self.state = match records_len {
None => FetchDecodeState::PartitionTags {
topics_left,
topic,
parts_left,
},
Some(blob_len) => FetchDecodeState::Blob {
topics_left,
topic,
parts_left,
partition: partition_index,
blob_len,
decode,
outcome_index,
},
};
}
FetchDecodeState::Blob {
topics_left,
topic,
parts_left,
partition,
blob_len,
decode,
outcome_index,
} => {
if buf.len().saturating_sub(pos) < blob_len {
self.state = FetchDecodeState::Blob {
topics_left,
topic,
parts_left,
partition,
blob_len,
decode,
outcome_index,
};
return Ok(pos);
}
if let Some(min_offset) = decode {
let blob = &buf[pos..pos + blob_len];
let stats = decode_record_batches(
&topic,
partition,
min_offset,
blob,
&mut self.builder,
)?;
if let Some(max_offset) = stats.max_offset {
self.outcomes[outcome_index].max_offset = Some(max_offset);
}
}
pos += blob_len;
self.state = FetchDecodeState::PartitionTags {
topics_left,
topic,
parts_left,
};
}
FetchDecodeState::PartitionTags {
topics_left,
topic,
parts_left,
} => match reader.skip_tags()? {
ChunkStep::Ready(()) => {
pos = reader.position();
self.state = FetchDecodeState::Partitions {
topics_left,
topic,
parts_left: parts_left - 1,
};
}
ChunkStep::NeedMore => {
self.state = FetchDecodeState::PartitionTags {
topics_left,
topic,
parts_left,
};
return Ok(pos);
}
},
FetchDecodeState::TopicTags { topics_left } => match reader.skip_tags()? {
ChunkStep::Ready(()) => {
pos = reader.position();
self.state = FetchDecodeState::Topics {
topics_left: topics_left - 1,
};
}
ChunkStep::NeedMore => {
self.state = FetchDecodeState::TopicTags { topics_left };
return Ok(pos);
}
},
FetchDecodeState::Trailer => match reader.skip_tags()? {
ChunkStep::Ready(()) => {
pos = reader.position();
self.state = FetchDecodeState::Done;
}
ChunkStep::NeedMore => {
self.state = FetchDecodeState::Trailer;
return Ok(pos);
}
},
FetchDecodeState::Done => {
if pos < buf.len() {
return Err(KafkaClientError::protocol(format!(
"Kafka Fetch response left {} trailing bytes",
buf.len() - pos
)));
}
self.state = FetchDecodeState::Done;
return Ok(pos);
}
}
}
}
}
enum ChunkStep<T> {
Ready(T),
NeedMore,
}
struct ChunkReader<'a> {
buf: &'a [u8],
pos: usize,
}
impl<'a> ChunkReader<'a> {
fn new(buf: &'a [u8], pos: usize) -> Self {
Self { buf, pos }
}
fn position(&self) -> usize {
self.pos
}
fn take(&mut self, len: usize) -> Option<&'a [u8]> {
let end = self.pos.checked_add(len)?;
if end > self.buf.len() {
return None;
}
let slice = &self.buf[self.pos..end];
self.pos = end;
Some(slice)
}
fn get_i16(&mut self) -> Option<i16> {
self.take(2)
.map(|bytes| i16::from_be_bytes(bytes.try_into().expect("2 bytes")))
}
fn get_i32(&mut self) -> Option<i32> {
self.take(4)
.map(|bytes| i32::from_be_bytes(bytes.try_into().expect("4 bytes")))
}
fn get_i64(&mut self) -> Option<i64> {
self.take(8)
.map(|bytes| i64::from_be_bytes(bytes.try_into().expect("8 bytes")))
}
fn get_unsigned_varint(&mut self) -> KafkaClientResult<Option<u32>> {
let mut value = 0_u32;
let mut cursor = self.pos;
for shift in (0..35).step_by(7) {
let Some(&byte) = self.buf.get(cursor) else {
return Ok(None);
};
cursor += 1;
value |= u32::from(byte & 0x7f) << shift;
if (byte & 0x80) == 0 {
self.pos = cursor;
return Ok(Some(value));
}
}
Err(KafkaClientError::protocol("Kafka unsigned varint overflow"))
}
fn get_compact_array_len(&mut self) -> KafkaClientResult<Option<usize>> {
match self.get_unsigned_varint()? {
None => Ok(None),
Some(0) => Err(KafkaClientError::protocol(
"non-null compact Kafka array had null length",
)),
Some(len) => Ok(Some((len - 1) as usize)),
}
}
fn get_compact_string(&mut self) -> KafkaClientResult<Option<String>> {
match self.get_unsigned_varint()? {
None => Ok(None),
Some(0) => Err(KafkaClientError::protocol(
"non-null compact Kafka string had null length",
)),
Some(len) => match self.take((len - 1) as usize) {
None => Ok(None),
Some(bytes) => std::str::from_utf8(bytes)
.map(|text| Some(text.to_owned()))
.map_err(|_| KafkaClientError::protocol("Kafka string was not valid UTF-8")),
},
}
}
fn get_compact_nullable_bytes_len(&mut self) -> KafkaClientResult<ChunkStep<Option<usize>>> {
match self.get_unsigned_varint()? {
None => Ok(ChunkStep::NeedMore),
Some(0) => Ok(ChunkStep::Ready(None)),
Some(len) => Ok(ChunkStep::Ready(Some((len - 1) as usize))),
}
}
fn skip_aborted_transactions(&mut self) -> KafkaClientResult<ChunkStep<()>> {
let count = match self.get_unsigned_varint()? {
None => return Ok(ChunkStep::NeedMore),
Some(0) => return Ok(ChunkStep::Ready(())),
Some(len) => len - 1,
};
for _ in 0..count {
let (Some(_producer_id), Some(_first_offset)) = (self.get_i64(), self.get_i64()) else {
return Ok(ChunkStep::NeedMore);
};
match self.skip_tags()? {
ChunkStep::Ready(()) => {}
ChunkStep::NeedMore => return Ok(ChunkStep::NeedMore),
}
}
Ok(ChunkStep::Ready(()))
}
fn skip_tags(&mut self) -> KafkaClientResult<ChunkStep<()>> {
let count = match self.get_unsigned_varint()? {
None => return Ok(ChunkStep::NeedMore),
Some(count) => count,
};
let mut previous: Option<u32> = None;
for _ in 0..count {
let Some(tag) = self.get_unsigned_varint()? else {
return Ok(ChunkStep::NeedMore);
};
if previous.is_some_and(|prev| tag <= prev) {
return Err(KafkaClientError::protocol(
"Kafka tagged fields were not strictly increasing",
));
}
previous = Some(tag);
let Some(size) = self.get_unsigned_varint()? else {
return Ok(ChunkStep::NeedMore);
};
if self.take(size as usize).is_none() {
return Ok(ChunkStep::NeedMore);
}
}
Ok(ChunkStep::Ready(()))
}
}
fn skip_i32_array(decoder: &mut Decoder<'_>, flexible: bool) -> KafkaClientResult<()> {
let len = decoder.get_array_len(flexible)?;
for _ in 0..len {
decoder.get_i32()?;
}
Ok(())
}
fn skip_aborted_transactions(decoder: &mut Decoder<'_>) -> KafkaClientResult<()> {
let Some(count) = decoder.get_nullable_array_len(true)? else {
return Ok(());
};
for _ in 0..count {
decoder.get_i64()?;
decoder.get_i64()?;
decoder.skip_tags()?;
}
Ok(())
}
fn expect_done(decoder: &Decoder<'_>, message: &'static str) -> KafkaClientResult<()> {
if decoder.is_done() {
Ok(())
} else {
Err(KafkaClientError::protocol(format!(
"{message} decode left {} trailing bytes at offset {}",
decoder.remaining(),
decoder.position()
)))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::native::{model::KafkaPayloadBatch, protocol::records::encode_test_record_batch};
type PartitionSpec = (i32, i16, i64, i64, Option<Vec<u8>>);
type BatchSummary = (Vec<(String, i32, i64, Vec<u8>)>, Vec<(String, i32, i64)>);
fn encode_fetch_body(top_error: i16, topics: &[(&str, Vec<PartitionSpec>)]) -> Vec<u8> {
let mut enc = Encoder::new();
enc.put_i32(0); enc.put_i16(top_error);
enc.put_i32(0); enc.put_array_len(topics.len(), true).expect("topics len");
for (topic, partitions) in topics {
enc.put_compact_string(topic).expect("topic");
enc.put_array_len(partitions.len(), true)
.expect("partitions len");
for (partition, error_code, high_watermark, log_start_offset, records) in partitions {
enc.put_i32(*partition);
enc.put_i16(*error_code);
enc.put_i64(*high_watermark);
enc.put_i64(*high_watermark); enc.put_i64(*log_start_offset);
enc.put_unsigned_varint(0); enc.put_i32(-1); match records {
None => enc.put_unsigned_varint(0), Some(blob) => {
enc.put_unsigned_varint((blob.len() + 1) as u32);
enc.put_raw(blob);
}
}
enc.put_empty_tags(); }
enc.put_empty_tags(); }
enc.put_empty_tags(); enc.into_inner()
}
fn oracle_decode(
body: &[u8],
offsets: &BTreeMap<(String, i32), i64>,
include_timestamps: bool,
) -> (Option<KafkaPayloadBatch>, i16, Vec<FetchPartitionOutcome>) {
let response = FetchResponse::decode(12, body).expect("oracle decode");
let mut builder = PayloadBatchBuilder::with_capacity(64, 4096, include_timestamps);
let mut outcomes = Vec::new();
for topic in &response.topics {
for partition in &topic.partitions {
let mut max_offset = None;
if partition.error_code == 0
&& let Some(&min) =
offsets.get(&(topic.topic.clone(), partition.partition_index))
&& let Some(records) = partition.records
{
let stats = decode_record_batches(
&topic.topic,
partition.partition_index,
min,
records,
&mut builder,
)
.expect("oracle record decode");
max_offset = stats.max_offset;
}
outcomes.push(FetchPartitionOutcome {
topic: topic.topic.clone(),
partition_index: partition.partition_index,
error_code: partition.error_code,
high_watermark: partition.high_watermark,
log_start_offset: partition.log_start_offset,
max_offset,
});
}
}
(builder.finish(Vec::new()), response.error_code, outcomes)
}
fn drive_streaming(
body: &[u8],
offsets: &BTreeMap<(String, i32), i64>,
include_timestamps: bool,
chunk: usize,
) -> KafkaClientResult<(Option<KafkaPayloadBatch>, i16, Vec<FetchPartitionOutcome>)> {
let mut decoder = FetchBodyDecoder::new(
12,
offsets.clone(),
PayloadBatchBuilder::with_capacity(64, 4096, include_timestamps),
)?;
let mut scratch = Vec::new();
let mut rest = body;
loop {
if !rest.is_empty() {
let take = chunk.max(1).min(rest.len());
scratch.extend_from_slice(&rest[..take]);
rest = &rest[take..];
}
let more = !rest.is_empty();
let consumed = decoder.feed(&scratch, more)?;
scratch.drain(0..consumed);
if !more {
decoder.finish()?;
break;
}
}
let (builder, error_code, outcomes) = decoder.into_parts();
Ok((builder.finish(Vec::new()), error_code, outcomes))
}
fn batch_summary(batch: &Option<KafkaPayloadBatch>) -> BatchSummary {
match batch {
None => (Vec::new(), Vec::new()),
Some(batch) => {
let records = batch
.records()
.iter()
.map(|record| {
(
record.topic(batch).to_owned(),
record.partition,
record.offset,
batch.payload(record).to_vec(),
)
})
.collect();
let watermarks = batch
.watermarks()
.iter()
.map(|commit| (commit.topic.clone(), commit.partition, commit.offset))
.collect();
(records, watermarks)
}
}
}
fn offsets(entries: &[(&str, i32, i64)]) -> BTreeMap<(String, i32), i64> {
entries
.iter()
.map(|(topic, partition, offset)| (((*topic).to_owned(), *partition), *offset))
.collect()
}
#[test]
fn streaming_fetch_decode_matches_oracle_across_every_chunk_boundary() {
let mut partial_tail = encode_test_record_batch(0, &[b"p0", b"p1"]).expect("partial batch");
partial_tail.extend_from_slice(&50_i64.to_be_bytes()); partial_tail.extend_from_slice(&999_i32.to_be_bytes()); partial_tail.extend_from_slice(&[7, 7, 7, 7]);
let batch_a1 = encode_test_record_batch(10, &[b"a1-0", b"a1-1", b"a1-2"]).expect("a1");
let batch_b0 = encode_test_record_batch(3, &[b"b0-0"]).expect("b0");
let batch_err = encode_test_record_batch(0, &[b"never-decoded"]).expect("err batch");
for include_timestamps in [false, true] {
let body = encode_fetch_body(
0,
&[
(
"topic-a",
vec![
(0, 0, 100, 0, Some(partial_tail.clone())),
(1, 0, 60, 5, Some(batch_a1.clone())),
],
),
(
"topic-b",
vec![
(0, 0, 20, 0, Some(batch_b0.clone())),
(2, 0, 0, 0, None), (3, 0, 0, 0, None), (4, 1, 0, 0, Some(batch_err.clone())), ],
),
],
);
let offsets = offsets(&[
("topic-a", 0, 0),
("topic-a", 1, 0),
("topic-b", 0, 0),
("topic-b", 2, 0),
("topic-b", 4, 0),
]);
let (oracle_batch, oracle_error, oracle_outcomes) =
oracle_decode(&body, &offsets, include_timestamps);
assert_eq!(oracle_error, 0);
let oracle_summary = batch_summary(&oracle_batch);
assert!(
!oracle_summary.0.is_empty(),
"oracle should decode some records"
);
for chunk in 1..=body.len() {
let (batch, error_code, outcomes) =
drive_streaming(&body, &offsets, include_timestamps, chunk)
.unwrap_or_else(|error| panic!("chunk {chunk}: {error}"));
assert_eq!(error_code, 0, "chunk {chunk}");
assert_eq!(outcomes, oracle_outcomes, "outcomes at chunk {chunk}");
assert_eq!(
batch_summary(&batch),
oracle_summary,
"records/watermarks at chunk {chunk}"
);
}
}
}
#[test]
fn streaming_fetch_decode_reports_response_level_error_code() {
let body = encode_fetch_body(7, &[]);
let offsets = BTreeMap::new();
for chunk in 1..=body.len() {
let (batch, error_code, outcomes) = drive_streaming(&body, &offsets, false, chunk)
.unwrap_or_else(|error| panic!("chunk {chunk}: {error}"));
assert_eq!(error_code, 7);
assert!(batch.is_none());
assert!(outcomes.is_empty());
}
}
#[test]
fn streaming_fetch_decode_surfaces_truncated_body() {
let batch = encode_test_record_batch(0, &[b"one", b"two"]).expect("batch");
let body = encode_fetch_body(0, &[("topic", vec![(0, 0, 10, 0, Some(batch))])]);
let offsets = offsets(&[("topic", 0, 0)]);
for cut in [1usize, 5, 12, body.len() / 2, body.len() - 1] {
let cut = cut.min(body.len() - 1).max(1);
let error = drive_streaming(&body[..cut], &offsets, false, 4)
.expect_err("truncated body must be a typed error");
assert!(
matches!(error, KafkaClientError::Protocol(_)),
"cut {cut}: {error:?}"
);
}
}
#[test]
fn streaming_fetch_decode_rejects_trailing_bytes() {
let batch = encode_test_record_batch(0, &[b"x"]).expect("batch");
let mut body = encode_fetch_body(0, &[("topic", vec![(0, 0, 4, 0, Some(batch))])]);
body.extend_from_slice(&[0xde, 0xad, 0xbe, 0xef]);
let offsets = offsets(&[("topic", 0, 0)]);
let error = drive_streaming(&body, &offsets, false, body.len())
.expect_err("trailing bytes must be a typed error");
assert!(matches!(error, KafkaClientError::Protocol(_)), "{error:?}");
}
#[test]
fn streaming_fetch_decode_surfaces_malformed_record_batch() {
let mut batch = encode_test_record_batch(0, &[b"corrupt-me"]).expect("batch");
let last = batch.len() - 1;
batch[last] ^= 0xff; let body = encode_fetch_body(0, &[("topic", vec![(0, 0, 4, 0, Some(batch))])]);
let offsets = offsets(&[("topic", 0, 0)]);
for chunk in [1usize, 7, body.len()] {
let error = drive_streaming(&body, &offsets, false, chunk)
.expect_err("bad CRC must surface a typed error");
assert!(matches!(error, KafkaClientError::Protocol(_)), "{error:?}");
}
}
#[test]
fn api_versions_response_round_trip_shape() {
let mut encoder = Encoder::new();
encoder.put_i16(0);
encoder.put_array_len(2, true).expect("array len");
for key in [1_i16, 18] {
encoder.put_i16(key);
encoder.put_i16(0);
encoder.put_i16(12);
encoder.put_empty_tags();
}
encoder.put_i32(0);
encoder.put_empty_tags();
let response =
ApiVersionsResponse::decode(3, &encoder.into_inner()).expect("decode api versions");
assert_eq!(response.error_code, 0);
assert_eq!(response.api_keys.len(), 2);
assert_eq!(response.api_keys[1].api_key, 18);
}
#[test]
fn metadata_request_encodes_flexible_tags() {
let bytes = encode_metadata_request(&["topic-a".to_owned()]).expect("encode metadata");
let mut decoder = Decoder::new(&bytes);
assert_eq!(decoder.get_array_len(true).expect("topics"), 1);
assert_eq!(
decoder.get_compact_string().expect("topic"),
"topic-a".to_owned()
);
decoder.skip_tags().expect("topic tags");
assert!(!decoder.get_bool().expect("auto create"));
assert!(!decoder.get_bool().expect("cluster auth"));
assert!(!decoder.get_bool().expect("topic auth"));
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
}
#[test]
fn fetch_response_rejects_short_records_length() {
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_i32(0);
encoder.put_array_len(1, true).expect("responses");
encoder.put_compact_string("topic").expect("topic");
encoder.put_array_len(1, true).expect("partitions");
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_i64(0);
encoder.put_i64(0);
encoder.put_i64(0);
encoder.put_unsigned_varint(0);
encoder.put_i32(-1);
encoder.put_unsigned_varint(4);
encoder.put_i8(1);
encoder.put_i8(2);
encoder.put_empty_tags();
encoder.put_empty_tags();
encoder.put_empty_tags();
let bytes = encoder.into_inner();
assert!(FetchResponse::decode(12, &bytes).is_err());
}
#[test]
fn produce_v9_request_and_response_round_trip_shape() {
let request = encode_produce_request(
None,
-1,
30_000,
&[ProduceTopicRequest {
name: "topic-a".to_owned(),
partitions: vec![ProducePartitionRequest {
index: 2,
records: vec![1, 2, 3],
}],
}],
)
.expect("produce request");
let mut decoder = Decoder::new(&request);
assert!(
decoder
.get_compact_nullable_string()
.expect("transactional id")
.is_none()
);
assert_eq!(decoder.get_i16().expect("acks"), -1);
assert_eq!(decoder.get_i32().expect("timeout"), 30_000);
assert_eq!(decoder.get_array_len(true).expect("topics"), 1);
assert_eq!(decoder.get_compact_string().expect("topic"), "topic-a");
assert_eq!(decoder.get_array_len(true).expect("partitions"), 1);
assert_eq!(decoder.get_i32().expect("partition"), 2);
assert_eq!(decoder.get_compact_bytes().expect("records"), [1, 2, 3]);
decoder.skip_tags().expect("partition tags");
decoder.skip_tags().expect("topic tags");
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
let mut encoder = Encoder::new();
encoder.put_array_len(1, true).expect("topics");
encoder.put_compact_string("topic-a").expect("topic");
encoder.put_array_len(1, true).expect("partitions");
encoder.put_i32(2);
encoder.put_i16(0);
encoder.put_i64(42);
encoder.put_i64(-1);
encoder.put_i64(0);
encoder.put_array_len(0, true).expect("record errors");
encoder
.put_compact_nullable_string(None)
.expect("error message");
encoder.put_empty_tags();
encoder.put_empty_tags();
encoder.put_i32(0);
encoder.put_empty_tags();
let results = decode_produce_response(9, &encoder.into_inner()).expect("response");
assert_eq!(results.len(), 1);
assert_eq!(results[0].topic, "topic-a");
assert_eq!(results[0].partition, 2);
assert_eq!(results[0].base_offset, 42);
}
#[test]
fn init_producer_id_v4_round_trip_shape() {
let request = encode_init_producer_id_request(None, -1, -1, -1)
.expect("encode init producer id request");
let mut decoder = Decoder::new(&request);
assert!(
decoder
.get_compact_nullable_string()
.expect("transactional id")
.is_none()
);
assert_eq!(decoder.get_i32().expect("transaction timeout"), -1);
assert_eq!(decoder.get_i64().expect("producer id"), -1);
assert_eq!(decoder.get_i16().expect("producer epoch"), -1);
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_i64(4242);
encoder.put_i16(7);
encoder.put_empty_tags();
let response = InitProducerIdResponse::decode(4, &encoder.into_inner())
.expect("decode init producer id response");
assert_eq!(response.error_code, 0);
assert_eq!(response.producer_id, 4242);
assert_eq!(response.producer_epoch, 7);
}
#[test]
fn init_producer_id_response_rejects_unimplemented_version() {
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_i64(1);
encoder.put_i16(0);
assert!(matches!(
InitProducerIdResponse::decode(3, &encoder.into_inner()),
Err(KafkaClientError::Unsupported(_))
));
}
#[test]
fn find_coordinator_v3_round_trip_shape() {
let request = encode_find_coordinator_request("group-a").expect("encode find coordinator");
let mut request_decoder = Decoder::new(&request);
assert_eq!(
request_decoder.get_compact_string().expect("group"),
"group-a"
);
assert_eq!(request_decoder.get_i8().expect("key type"), 0);
request_decoder.skip_tags().expect("request tags");
assert!(request_decoder.is_done());
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_compact_nullable_string(None).expect("message");
encoder.put_i32(1);
encoder.put_compact_string("localhost").expect("host");
encoder.put_i32(9092);
encoder.put_empty_tags();
let response =
FindCoordinatorResponse::decode(3, &encoder.into_inner()).expect("decode coordinator");
assert_eq!(response.node_id, 1);
assert_eq!(response.host, "localhost");
assert_eq!(response.port, 9092);
}
#[test]
fn offset_commit_v8_round_trip_shape() {
let request = encode_offset_commit_request(
"group-a",
-1,
"",
None,
&[OffsetCommitTopicRequest {
name: "topic-a".to_owned(),
partitions: vec![OffsetCommitPartitionRequest {
partition_index: 2,
committed_offset: 42,
committed_leader_epoch: -1,
committed_metadata: None,
}],
}],
)
.expect("encode offset commit");
let mut request_decoder = Decoder::new(&request);
assert_eq!(
request_decoder.get_compact_string().expect("group"),
"group-a"
);
assert_eq!(request_decoder.get_i32().expect("generation"), -1);
assert_eq!(request_decoder.get_compact_string().expect("member"), "");
assert!(
request_decoder
.get_compact_nullable_string()
.expect("instance")
.is_none()
);
assert_eq!(request_decoder.get_array_len(true).expect("topics"), 1);
assert_eq!(
request_decoder.get_compact_string().expect("topic"),
"topic-a"
);
assert_eq!(request_decoder.get_array_len(true).expect("partitions"), 1);
assert_eq!(request_decoder.get_i32().expect("partition"), 2);
assert_eq!(request_decoder.get_i64().expect("offset"), 42);
assert_eq!(request_decoder.get_i32().expect("leader epoch"), -1);
assert!(
request_decoder
.get_compact_nullable_string()
.expect("metadata")
.is_none()
);
request_decoder.skip_tags().expect("partition tags");
request_decoder.skip_tags().expect("topic tags");
request_decoder.skip_tags().expect("request tags");
assert!(request_decoder.is_done());
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_array_len(1, true).expect("topics");
encoder.put_compact_string("topic-a").expect("topic");
encoder.put_array_len(1, true).expect("partitions");
encoder.put_i32(2);
encoder.put_i16(0);
encoder.put_empty_tags();
encoder.put_empty_tags();
encoder.put_empty_tags();
let results =
decode_offset_commit_response(8, &encoder.into_inner()).expect("commit response");
assert_eq!(results.len(), 1);
assert_eq!(results[0].partition_index, 2);
assert_eq!(results[0].error_code, 0);
}
#[test]
fn offset_fetch_v7_round_trip_shape() {
let request = encode_offset_fetch_request(
"group-a",
&[OffsetFetchTopicRequest {
name: "topic-a".to_owned(),
partitions: vec![0, 2],
}],
)
.expect("encode offset fetch");
let mut request_decoder = Decoder::new(&request);
assert_eq!(
request_decoder.get_compact_string().expect("group"),
"group-a"
);
assert_eq!(request_decoder.get_array_len(true).expect("topics"), 1);
assert_eq!(
request_decoder.get_compact_string().expect("topic"),
"topic-a"
);
assert_eq!(request_decoder.get_array_len(true).expect("partitions"), 2);
assert_eq!(request_decoder.get_i32().expect("partition 0"), 0);
assert_eq!(request_decoder.get_i32().expect("partition 2"), 2);
request_decoder.skip_tags().expect("topic tags");
assert!(!request_decoder.get_bool().expect("require stable"));
request_decoder.skip_tags().expect("request tags");
assert!(request_decoder.is_done());
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_array_len(1, true).expect("topics");
encoder.put_compact_string("topic-a").expect("topic");
encoder.put_array_len(1, true).expect("partitions");
encoder.put_i32(2);
encoder.put_i64(42);
encoder.put_i32(-1);
encoder.put_compact_nullable_string(None).expect("metadata");
encoder.put_i16(0);
encoder.put_empty_tags();
encoder.put_empty_tags();
encoder.put_i16(0);
encoder.put_empty_tags();
let response =
OffsetFetchResponse::decode(7, &encoder.into_inner()).expect("offset fetch response");
assert_eq!(response.error_code, 0);
assert_eq!(response.offsets[0].committed_offset, 42);
}
#[test]
fn group_requests_round_trip_shape() {
let metadata = encode_consumer_protocol_subscription(
&["topic-a".to_owned()],
&[TopicPartition::new("topic-a", 1)],
7,
)
.expect("subscription");
let subscription =
decode_consumer_protocol_subscription(&metadata).expect("decode subscription");
assert_eq!(subscription.topics, vec!["topic-a".to_owned()]);
assert_eq!(
subscription.owned_partitions,
vec![TopicPartition::new("topic-a", 1)]
);
assert_eq!(subscription.generation_id, 7);
let request = encode_join_group_request(
"group-a",
45_000,
300_000,
"member-a",
Some("instance-a"),
&[ConsumerGroupProtocol {
name: "cooperative-sticky".to_owned(),
metadata: metadata.clone(),
}],
)
.expect("join request");
let mut decoder = Decoder::new(&request);
assert_eq!(decoder.get_compact_string().expect("group"), "group-a");
assert_eq!(decoder.get_i32().expect("session"), 45_000);
assert_eq!(decoder.get_i32().expect("rebalance"), 300_000);
assert_eq!(decoder.get_compact_string().expect("member"), "member-a");
assert_eq!(
decoder
.get_compact_nullable_string()
.expect("instance")
.as_deref(),
Some("instance-a")
);
assert_eq!(
decoder.get_compact_string().expect("protocol type"),
"consumer"
);
assert_eq!(decoder.get_array_len(true).expect("protocols"), 1);
assert_eq!(
decoder.get_compact_string().expect("protocol"),
"cooperative-sticky"
);
assert_eq!(decoder.get_compact_bytes().expect("metadata"), metadata);
decoder.skip_tags().expect("protocol tags");
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
let assignment = encode_consumer_protocol_assignment(&[TopicPartition::new("topic-a", 2)])
.expect("assignment");
assert_eq!(
decode_consumer_protocol_assignment(&assignment).expect("decode assignment"),
vec![TopicPartition::new("topic-a", 2)]
);
let request = encode_sync_group_request(
"group-a",
7,
"member-a",
Some("instance-a"),
&[SyncGroupAssignment {
member_id: "member-a".to_owned(),
assignment: assignment.clone(),
}],
)
.expect("sync request");
let mut decoder = Decoder::new(&request);
assert_eq!(decoder.get_compact_string().expect("group"), "group-a");
assert_eq!(decoder.get_i32().expect("generation"), 7);
assert_eq!(decoder.get_compact_string().expect("member"), "member-a");
assert_eq!(
decoder
.get_compact_nullable_string()
.expect("instance")
.as_deref(),
Some("instance-a")
);
assert_eq!(decoder.get_array_len(true).expect("assignments"), 1);
assert_eq!(
decoder.get_compact_string().expect("assignment member"),
"member-a"
);
assert_eq!(decoder.get_compact_bytes().expect("assignment"), assignment);
decoder.skip_tags().expect("assignment tags");
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
let heartbeat = encode_heartbeat_request("group-a", 7, "member-a", Some("instance-a"))
.expect("heartbeat");
let mut decoder = Decoder::new(&heartbeat);
assert_eq!(decoder.get_compact_string().expect("group"), "group-a");
assert_eq!(decoder.get_i32().expect("generation"), 7);
assert_eq!(decoder.get_compact_string().expect("member"), "member-a");
assert_eq!(
decoder
.get_compact_nullable_string()
.expect("instance")
.as_deref(),
Some("instance-a")
);
decoder.skip_tags().expect("heartbeat tags");
assert!(decoder.is_done());
let leave =
encode_leave_group_request("group-a", "member-a", Some("instance-a")).expect("leave");
let mut decoder = Decoder::new(&leave);
assert_eq!(decoder.get_compact_string().expect("group"), "group-a");
assert_eq!(decoder.get_array_len(true).expect("members"), 1);
assert_eq!(decoder.get_compact_string().expect("member"), "member-a");
assert_eq!(
decoder
.get_compact_nullable_string()
.expect("instance")
.as_deref(),
Some("instance-a")
);
decoder.skip_tags().expect("member tags");
decoder.skip_tags().expect("request tags");
assert!(decoder.is_done());
}
#[test]
fn group_responses_decode_shape() {
let metadata = encode_consumer_protocol_subscription(&["topic-a".to_owned()], &[], -1)
.expect("subscription");
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_i32(3);
encoder
.put_compact_string("cooperative-sticky")
.expect("protocol");
encoder.put_compact_string("member-a").expect("leader");
encoder.put_compact_string("member-a").expect("member");
encoder.put_array_len(1, true).expect("members");
encoder.put_compact_string("member-a").expect("member id");
encoder
.put_compact_nullable_string(Some("instance-a"))
.expect("instance");
encoder.put_compact_bytes(&metadata).expect("metadata");
encoder.put_empty_tags();
encoder.put_empty_tags();
let join = JoinGroupResponse::decode(6, &encoder.into_inner()).expect("join decode");
assert_eq!(join.generation_id, 3);
assert_eq!(join.members[0].metadata, metadata);
let assignment = encode_consumer_protocol_assignment(&[TopicPartition::new("topic-a", 0)])
.expect("assignment");
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder
.put_compact_bytes(&assignment)
.expect("assignment bytes");
encoder.put_empty_tags();
let sync = SyncGroupResponse::decode(4, &encoder.into_inner()).expect("sync decode");
assert_eq!(sync.assignment, assignment);
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(27);
encoder.put_empty_tags();
let heartbeat =
HeartbeatResponse::decode(4, &encoder.into_inner()).expect("heartbeat decode");
assert_eq!(heartbeat.error_code, 27);
let mut encoder = Encoder::new();
encoder.put_i32(0);
encoder.put_i16(0);
encoder.put_array_len(1, true).expect("members");
encoder.put_compact_string("member-a").expect("member");
encoder
.put_compact_nullable_string(Some("instance-a"))
.expect("instance");
encoder.put_i16(0);
encoder.put_empty_tags();
encoder.put_empty_tags();
let leave = LeaveGroupResponse::decode(4, &encoder.into_inner()).expect("leave decode");
assert_eq!(leave.members[0].member_id, "member-a");
}
}