use bytes::{Buf, BufMut};
use super::{VersionedDecode, VersionedEncode, non_nullable_string};
use crate::error::{ErrorCode, ProtocolErrorKind, Result};
use crate::protocol::api::ApiKey;
use crate::protocol::primitives::{Decode, Encode, KafkaString, TaggedFields, TryEncode};
use crate::protocol::{check_compact_array_len, decode_capacity, encode_compact_array_len};
use crate::util::varint::decode_unsigned_varint;
#[derive(Debug, Clone)]
pub struct DescribeQuorumPartitionRequest {
pub partition_index: i32,
}
#[derive(Debug, Clone)]
pub struct DescribeQuorumTopicRequest {
pub topic_name: String,
pub partitions: Vec<DescribeQuorumPartitionRequest>,
}
#[derive(Debug, Clone)]
pub struct DescribeQuorumRequest {
pub topics: Vec<DescribeQuorumTopicRequest>,
}
impl DescribeQuorumRequest {
pub fn api_key() -> ApiKey {
ApiKey::DescribeQuorum
}
pub fn encode_v0(&self, buf: &mut impl BufMut) -> Result<()> {
encode_compact_array_len(self.topics.len(), buf)?;
for topic in &self.topics {
KafkaString::new(&topic.topic_name).try_encode_compact(buf)?;
encode_compact_array_len(topic.partitions.len(), buf)?;
for partition in &topic.partitions {
partition.partition_index.encode(buf);
TaggedFields::default().try_encode(buf)?;
}
TaggedFields::default().try_encode(buf)?;
}
TaggedFields::default().try_encode(buf)?;
Ok(())
}
}
impl VersionedEncode for DescribeQuorumRequest {
fn encode_versioned(&self, version: i16, buf: &mut impl BufMut) -> Result<()> {
match version {
0..=2 => self.encode_v0(buf),
_ => unsupported_encode!("DescribeQuorumRequest", version),
}
}
}
#[derive(Debug, Clone)]
pub struct QuorumReplicaState {
pub replica_id: i32,
pub log_end_offset: i64,
pub last_fetch_timestamp: i64,
pub last_caught_up_timestamp: i64,
pub replica_directory_id: Option<[u8; 16]>,
}
#[derive(Debug, Clone)]
pub struct QuorumListener {
pub name: String,
pub host: String,
pub port: u16,
}
#[derive(Debug, Clone)]
pub struct QuorumNode {
pub node_id: i32,
pub listeners: Vec<QuorumListener>,
}
#[derive(Debug, Clone)]
pub struct DescribeQuorumPartitionResponse {
pub partition_index: i32,
pub error_code: ErrorCode,
pub error_message: Option<String>,
pub leader_id: i32,
pub leader_epoch: i32,
pub high_watermark: i64,
pub current_voters: Vec<QuorumReplicaState>,
pub observers: Vec<QuorumReplicaState>,
}
#[derive(Debug, Clone)]
pub struct DescribeQuorumTopicResponse {
pub topic_name: String,
pub partitions: Vec<DescribeQuorumPartitionResponse>,
}
#[derive(Debug, Clone)]
pub struct DescribeQuorumResponse {
pub error_code: ErrorCode,
pub error_message: Option<String>,
pub topics: Vec<DescribeQuorumTopicResponse>,
pub nodes: Vec<QuorumNode>,
}
impl DescribeQuorumResponse {
fn decode_replica_states(buf: &mut impl Buf, version: i16) -> Result<Vec<QuorumReplicaState>> {
let count = check_compact_array_len(decode_unsigned_varint(buf)?)? as usize;
let mut states = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let replica_id = i32::decode(buf)?;
let replica_directory_id = if version >= 2 {
Some(read_uuid(buf)?)
} else {
None
};
let log_end_offset = i64::decode(buf)?;
let (last_fetch_timestamp, last_caught_up_timestamp) = if version >= 1 {
(i64::decode(buf)?, i64::decode(buf)?)
} else {
(-1, -1)
};
TaggedFields::decode(buf)?;
states.push(QuorumReplicaState {
replica_id,
log_end_offset,
last_fetch_timestamp,
last_caught_up_timestamp,
replica_directory_id,
});
}
Ok(states)
}
fn decode_nodes(buf: &mut impl Buf) -> Result<Vec<QuorumNode>> {
let count = check_compact_array_len(decode_unsigned_varint(buf)?)? as usize;
let mut nodes = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let node_id = i32::decode(buf)?;
let listener_count = check_compact_array_len(decode_unsigned_varint(buf)?)? as usize;
let mut listeners =
Vec::with_capacity(decode_capacity(listener_count, buf.remaining()));
for _ in 0..listener_count {
let name =
non_nullable_string("listener name", KafkaString::decode_compact(buf)?.0)?;
let host =
non_nullable_string("listener host", KafkaString::decode_compact(buf)?.0)?;
if buf.remaining() < 2 {
return Err(crate::error::KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
"not enough bytes for listener port",
));
}
let port = buf.get_u16();
TaggedFields::decode(buf)?;
listeners.push(QuorumListener { name, host, port });
}
TaggedFields::decode(buf)?;
nodes.push(QuorumNode { node_id, listeners });
}
Ok(nodes)
}
pub fn decode_v0(buf: &mut impl Buf) -> Result<Self> {
Self::decode_inner(buf, 0)
}
pub fn decode_v1(buf: &mut impl Buf) -> Result<Self> {
Self::decode_inner(buf, 1)
}
pub fn decode_v2(buf: &mut impl Buf) -> Result<Self> {
Self::decode_inner(buf, 2)
}
fn decode_inner(buf: &mut impl Buf, version: i16) -> Result<Self> {
let error_code = ErrorCode::from(i16::decode(buf)?);
let error_message = if version >= 2 {
KafkaString::decode_compact(buf)?.0
} else {
None
};
let topic_count = check_compact_array_len(decode_unsigned_varint(buf)?)? as usize;
let mut topics = Vec::with_capacity(decode_capacity(topic_count, buf.remaining()));
for _ in 0..topic_count {
let topic_name = {
let len = decode_unsigned_varint(buf)? as usize;
if len < 1 {
return Err(crate::error::KrafkaError::protocol_kind(
ProtocolErrorKind::Malformed,
"compact string length 0 is null but field is non-nullable",
));
}
let str_len = len - 1;
if buf.remaining() < str_len {
return Err(crate::error::KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
"not enough bytes for compact string",
));
}
let bytes = buf.copy_to_bytes(str_len);
String::from_utf8(bytes.to_vec()).map_err(|e| {
crate::error::KrafkaError::protocol_kind(
ProtocolErrorKind::InvalidUtf8,
format!("invalid UTF-8: {e}"),
)
})?
};
let partition_count = check_compact_array_len(decode_unsigned_varint(buf)?)? as usize;
let mut partitions =
Vec::with_capacity(decode_capacity(partition_count, buf.remaining()));
for _ in 0..partition_count {
let partition_index = i32::decode(buf)?;
let partition_error_code = ErrorCode::from(i16::decode(buf)?);
let partition_error_message = if version >= 2 {
KafkaString::decode_compact(buf)?.0
} else {
None
};
let leader_id = i32::decode(buf)?;
let leader_epoch = i32::decode(buf)?;
let high_watermark = i64::decode(buf)?;
let current_voters = Self::decode_replica_states(buf, version)?;
let observers = Self::decode_replica_states(buf, version)?;
TaggedFields::decode(buf)?;
partitions.push(DescribeQuorumPartitionResponse {
partition_index,
error_code: partition_error_code,
error_message: partition_error_message,
leader_id,
leader_epoch,
high_watermark,
current_voters,
observers,
});
}
TaggedFields::decode(buf)?;
topics.push(DescribeQuorumTopicResponse {
topic_name,
partitions,
});
}
let nodes = if version >= 2 {
Self::decode_nodes(buf)?
} else {
Vec::new()
};
TaggedFields::decode(buf)?;
Ok(Self {
error_code,
error_message,
topics,
nodes,
})
}
}
fn read_uuid(buf: &mut impl Buf) -> Result<[u8; 16]> {
if buf.remaining() < 16 {
return Err(crate::error::KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
"not enough bytes for replica_directory_id UUID",
));
}
let mut id = [0u8; 16];
buf.copy_to_slice(&mut id);
Ok(id)
}
impl VersionedDecode for DescribeQuorumResponse {
fn decode_versioned(version: i16, buf: &mut impl Buf) -> Result<Self> {
match version {
0 => Self::decode_v0(buf),
1 => Self::decode_v1(buf),
2 => Self::decode_v2(buf),
_ => unsupported_decode!("DescribeQuorumResponse", version),
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use bytes::BytesMut;
#[test]
fn describe_quorum_request_encode_v0() {
let request = DescribeQuorumRequest {
topics: vec![DescribeQuorumTopicRequest {
topic_name: "__cluster_metadata".to_string(),
partitions: vec![DescribeQuorumPartitionRequest { partition_index: 0 }],
}],
};
let mut buf = BytesMut::new();
request.encode_v0(&mut buf).unwrap();
assert!(!buf.is_empty());
}
#[test]
fn describe_quorum_response_decode_v0() {
let mut buf = BytesMut::new();
buf.put_i16(0);
buf.put_u8(2);
let name = b"__cluster_metadata";
buf.put_u8((name.len() + 1) as u8);
buf.put_slice(name);
buf.put_u8(2);
buf.put_i32(0);
buf.put_i16(0);
buf.put_i32(1);
buf.put_i32(5);
buf.put_i64(100);
buf.put_u8(3);
buf.put_i32(1);
buf.put_i64(100);
buf.put_u8(0); buf.put_i32(2);
buf.put_i64(98);
buf.put_u8(0); buf.put_u8(2);
buf.put_i32(3);
buf.put_i64(95);
buf.put_u8(0); buf.put_u8(0);
buf.put_u8(0);
buf.put_u8(0);
let mut read_buf = buf.freeze();
let response = DescribeQuorumResponse::decode_v0(&mut read_buf).unwrap();
assert!(response.error_code.is_ok());
assert_eq!(response.topics.len(), 1);
assert_eq!(response.topics[0].topic_name, "__cluster_metadata");
let partition = &response.topics[0].partitions[0];
assert_eq!(partition.partition_index, 0);
assert!(partition.error_code.is_ok());
assert_eq!(partition.leader_id, 1);
assert_eq!(partition.leader_epoch, 5);
assert_eq!(partition.high_watermark, 100);
assert_eq!(partition.current_voters.len(), 2);
assert_eq!(partition.current_voters[0].replica_id, 1);
assert_eq!(partition.current_voters[0].log_end_offset, 100);
assert_eq!(partition.current_voters[1].replica_id, 2);
assert_eq!(partition.current_voters[1].log_end_offset, 98);
assert_eq!(partition.observers.len(), 1);
assert_eq!(partition.observers[0].replica_id, 3);
assert_eq!(partition.observers[0].log_end_offset, 95);
}
#[test]
fn describe_quorum_request_roundtrip_v0() {
let request = DescribeQuorumRequest {
topics: vec![DescribeQuorumTopicRequest {
topic_name: "__cluster_metadata".to_string(),
partitions: vec![
DescribeQuorumPartitionRequest { partition_index: 0 },
DescribeQuorumPartitionRequest { partition_index: 1 },
],
}],
};
let mut buf = BytesMut::new();
request.encode_versioned(0, &mut buf).unwrap();
assert!(!buf.is_empty());
}
#[test]
fn describe_quorum_response_empty_topics() {
let mut buf = BytesMut::new();
buf.put_i16(0);
buf.put_u8(1);
buf.put_u8(0);
let mut read_buf = buf.freeze();
let response = DescribeQuorumResponse::decode_v0(&mut read_buf).unwrap();
assert!(response.topics.is_empty());
}
#[test]
fn describe_quorum_versioned_encode_dispatch() {
let request = DescribeQuorumRequest { topics: vec![] };
let mut buf = BytesMut::new();
request.encode_versioned(0, &mut buf).unwrap();
let mut buf2 = BytesMut::new();
assert!(request.encode_versioned(99, &mut buf2).is_err());
}
#[test]
fn describe_quorum_versioned_decode_dispatch() {
let mut buf = BytesMut::new();
buf.put_i16(0); buf.put_u8(1); buf.put_u8(0);
let mut read_buf = buf.freeze();
DescribeQuorumResponse::decode_versioned(0, &mut read_buf).unwrap();
let mut empty = BytesMut::new().freeze();
assert!(DescribeQuorumResponse::decode_versioned(99, &mut empty).is_err());
}
#[test]
fn describe_quorum_response_with_error() {
let mut buf = BytesMut::new();
buf.put_i16(-1);
buf.put_u8(1);
buf.put_u8(0);
let mut read_buf = buf.freeze();
let response = DescribeQuorumResponse::decode_v0(&mut read_buf).unwrap();
assert!(!response.error_code.is_ok());
}
#[test]
fn describe_quorum_response_decode_v1_timestamps() {
let mut buf = BytesMut::new();
buf.put_i16(0); buf.put_u8(2); let name = b"__cluster_metadata";
buf.put_u8((name.len() + 1) as u8);
buf.put_slice(name);
buf.put_u8(2); buf.put_i32(0); buf.put_i16(0); buf.put_i32(1); buf.put_i32(5); buf.put_i64(100); buf.put_u8(3);
buf.put_i32(1);
buf.put_i64(100);
buf.put_i64(-1); buf.put_i64(-1);
buf.put_u8(0);
buf.put_i32(2);
buf.put_i64(98);
buf.put_i64(1_700_000_000_000);
buf.put_i64(1_699_999_999_000);
buf.put_u8(0);
buf.put_u8(2);
buf.put_i32(3);
buf.put_i64(95);
buf.put_i64(1_700_000_000_500);
buf.put_i64(1_699_999_998_000);
buf.put_u8(0);
buf.put_u8(0); buf.put_u8(0); buf.put_u8(0);
let response = DescribeQuorumResponse::decode_versioned(1, &mut buf.freeze())
.expect("v1 response decodes");
let partition = &response.topics[0].partitions[0];
assert_eq!(partition.current_voters.len(), 2);
assert_eq!(
partition.current_voters[1].last_fetch_timestamp,
1_700_000_000_000
);
assert_eq!(
partition.current_voters[1].last_caught_up_timestamp,
1_699_999_999_000
);
assert_eq!(partition.observers.len(), 1);
assert_eq!(
partition.observers[0].last_fetch_timestamp,
1_700_000_000_500
);
assert!(response.error_code.is_ok());
}
#[test]
fn describe_quorum_response_v0_timestamps_are_unknown() {
let mut buf = BytesMut::new();
buf.put_i16(0);
buf.put_u8(2);
let name = b"t";
buf.put_u8((name.len() + 1) as u8);
buf.put_slice(name);
buf.put_u8(2);
buf.put_i32(0);
buf.put_i16(0);
buf.put_i32(1);
buf.put_i32(0);
buf.put_i64(0);
buf.put_u8(2); buf.put_i32(1);
buf.put_i64(10);
buf.put_u8(0);
buf.put_u8(1); buf.put_u8(0);
buf.put_u8(0);
buf.put_u8(0);
let response = DescribeQuorumResponse::decode_v0(&mut buf.freeze()).unwrap();
let voter = &response.topics[0].partitions[0].current_voters[0];
assert_eq!(voter.last_fetch_timestamp, -1);
assert_eq!(voter.last_caught_up_timestamp, -1);
}
#[test]
fn describe_quorum_request_v1_matches_v0() {
let request = DescribeQuorumRequest {
topics: vec![DescribeQuorumTopicRequest {
topic_name: "__cluster_metadata".to_string(),
partitions: vec![DescribeQuorumPartitionRequest { partition_index: 0 }],
}],
};
let mut v0 = BytesMut::new();
let mut v1 = BytesMut::new();
request.encode_versioned(0, &mut v0).unwrap();
request.encode_versioned(1, &mut v1).unwrap();
assert_eq!(v0, v1);
}
#[test]
fn describe_quorum_response_decode_v2_full_body() {
let mut buf = BytesMut::new();
buf.put_i16(0); put_compact_null_string(&mut buf); buf.put_u8(2); put_compact_string(&mut buf, "__cluster_metadata");
buf.put_u8(2); buf.put_i32(0); buf.put_i16(0); put_compact_null_string(&mut buf); buf.put_i32(1); buf.put_i32(5); buf.put_i64(100); buf.put_u8(3);
buf.put_i32(1);
buf.put_slice(&[0xAA; 16]); buf.put_i64(100);
buf.put_i64(-1);
buf.put_i64(-1);
buf.put_u8(0);
buf.put_i32(2);
buf.put_slice(&[0xBB; 16]);
buf.put_i64(98);
buf.put_i64(1_700_000_000_000);
buf.put_i64(1_699_999_999_000);
buf.put_u8(0);
buf.put_u8(2);
buf.put_i32(3);
buf.put_slice(&[0xCC; 16]);
buf.put_i64(95);
buf.put_i64(1_700_000_000_500);
buf.put_i64(1_699_999_998_000);
buf.put_u8(0);
buf.put_u8(0); buf.put_u8(0); buf.put_u8(3);
buf.put_i32(1);
buf.put_u8(2); put_compact_string(&mut buf, "CONTROLLER");
put_compact_string(&mut buf, "broker-1.internal");
buf.put_u16(9093);
buf.put_u8(0); buf.put_u8(0); buf.put_i32(2);
buf.put_u8(1); buf.put_u8(0); buf.put_u8(0);
let response = DescribeQuorumResponse::decode_versioned(2, &mut buf.freeze())
.expect("v2 response decodes");
let partition = &response.topics[0].partitions[0];
assert_eq!(
partition.leader_id, 1,
"leader_id must survive the v2 inserts"
);
assert_eq!(partition.high_watermark, 100);
assert_eq!(partition.current_voters.len(), 2);
assert_eq!(
partition.current_voters[0].replica_directory_id,
Some([0xAA; 16])
);
assert_eq!(partition.current_voters[1].log_end_offset, 98);
assert_eq!(
partition.current_voters[1].last_fetch_timestamp,
1_700_000_000_000
);
assert_eq!(
partition.observers[0].replica_directory_id,
Some([0xCC; 16])
);
assert_eq!(response.nodes.len(), 2);
assert_eq!(response.nodes[0].node_id, 1);
assert_eq!(response.nodes[0].listeners.len(), 1);
assert_eq!(response.nodes[0].listeners[0].name, "CONTROLLER");
assert_eq!(response.nodes[0].listeners[0].host, "broker-1.internal");
assert_eq!(response.nodes[0].listeners[0].port, 9093);
assert!(response.nodes[1].listeners.is_empty());
}
#[test]
fn describe_quorum_v2_port_is_unsigned() {
let mut buf = BytesMut::new();
buf.put_i16(0);
put_compact_null_string(&mut buf);
buf.put_u8(1); buf.put_u8(2); buf.put_i32(7);
buf.put_u8(2); put_compact_string(&mut buf, "CONTROLLER");
put_compact_string(&mut buf, "h");
buf.put_u16(49152);
buf.put_u8(0);
buf.put_u8(0);
buf.put_u8(0);
let response = DescribeQuorumResponse::decode_versioned(2, &mut buf.freeze()).unwrap();
assert_eq!(response.nodes[0].listeners[0].port, 49152);
}
#[test]
fn describe_quorum_v1_has_no_v2_fields() {
let mut buf = BytesMut::new();
buf.put_i16(0);
buf.put_u8(2);
put_compact_string(&mut buf, "t");
buf.put_u8(2);
buf.put_i32(0);
buf.put_i16(0);
buf.put_i32(1);
buf.put_i32(0);
buf.put_i64(0);
buf.put_u8(2); buf.put_i32(1);
buf.put_i64(10);
buf.put_i64(-1);
buf.put_i64(-1);
buf.put_u8(0);
buf.put_u8(1); buf.put_u8(0);
buf.put_u8(0);
buf.put_u8(0);
let response = DescribeQuorumResponse::decode_versioned(1, &mut buf.freeze()).unwrap();
assert!(response.error_message.is_none());
assert!(response.nodes.is_empty());
let partition = &response.topics[0].partitions[0];
assert!(partition.error_message.is_none());
assert!(partition.current_voters[0].replica_directory_id.is_none());
}
#[test]
fn describe_quorum_request_v2_matches_v0() {
let request = DescribeQuorumRequest {
topics: vec![DescribeQuorumTopicRequest {
topic_name: "__cluster_metadata".to_string(),
partitions: vec![DescribeQuorumPartitionRequest { partition_index: 0 }],
}],
};
let mut v0 = BytesMut::new();
let mut v2 = BytesMut::new();
request.encode_versioned(0, &mut v0).unwrap();
request.encode_versioned(2, &mut v2).unwrap();
assert_eq!(v0, v2);
let mut v3 = BytesMut::new();
assert!(request.encode_versioned(3, &mut v3).is_err());
}
#[test]
fn describe_quorum_v2_truncated_directory_id_errors() {
let mut buf = BytesMut::new();
buf.put_i16(0);
put_compact_null_string(&mut buf);
buf.put_u8(2);
put_compact_string(&mut buf, "t");
buf.put_u8(2);
buf.put_i32(0);
buf.put_i16(0);
put_compact_null_string(&mut buf);
buf.put_i32(1);
buf.put_i32(0);
buf.put_i64(0);
buf.put_u8(2); buf.put_i32(1);
buf.put_slice(&[0xAA; 8]);
let err = DescribeQuorumResponse::decode_versioned(2, &mut buf.freeze())
.expect_err("short UUID must error");
assert!(
err.to_string().contains("replica_directory_id"),
"got: {err}"
);
}
fn put_compact_string(buf: &mut BytesMut, s: &str) {
crate::util::varint::encode_unsigned_varint((s.len() + 1) as u32, buf);
buf.put_slice(s.as_bytes());
}
fn put_compact_null_string(buf: &mut BytesMut) {
crate::util::varint::encode_unsigned_varint(0, buf);
}
#[test]
fn describe_quorum_response_empty_voters_and_observers() {
let mut buf = BytesMut::new();
buf.put_i16(0); buf.put_u8(2); let name = b"t";
buf.put_u8((name.len() + 1) as u8);
buf.put_slice(name);
buf.put_u8(2); buf.put_i32(0); buf.put_i16(0); buf.put_i32(1); buf.put_i32(0); buf.put_i64(0); buf.put_u8(1); buf.put_u8(1); buf.put_u8(0); buf.put_u8(0); buf.put_u8(0);
let mut read_buf = buf.freeze();
let response = DescribeQuorumResponse::decode_v0(&mut read_buf).unwrap();
let partition = &response.topics[0].partitions[0];
assert!(partition.current_voters.is_empty());
assert!(partition.observers.is_empty());
}
}