use std::collections::{BTreeMap, BTreeSet};
use crate::native::{
KafkaClientError, KafkaClientResult,
model::TopicPartition,
protocol::{Decoder, Encoder},
};
#[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, 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 })
}
}
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::*;
#[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 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");
}
}