use crate::protocol::decode_capacity;
use bytes::{Buf, BufMut};
use super::{ConfigResourceType, VersionedDecode, VersionedEncode, non_nullable_string};
use crate::error::{ErrorCode, Result};
use crate::protocol::check_compact_array_len;
use crate::protocol::encode_compact_array_len;
use crate::protocol::primitives::{Decode, Encode, KafkaString, TaggedFields, TryEncode};
#[derive(Debug, Clone, Default)]
pub struct ListConfigResourcesRequest {
pub resource_types: Vec<ConfigResourceType>,
}
impl ListConfigResourcesRequest {
pub fn all() -> Self {
Self::default()
}
pub fn with_types(resource_types: Vec<ConfigResourceType>) -> Self {
Self { resource_types }
}
pub fn encode_v0(&self, buf: &mut impl BufMut) -> Result<()> {
TaggedFields::default().try_encode(buf)?;
Ok(())
}
pub fn encode_v1(&self, buf: &mut impl BufMut) -> Result<()> {
encode_compact_array_len(self.resource_types.len(), buf)?;
for ty in &self.resource_types {
ty.to_i8().encode(buf);
}
TaggedFields::default().try_encode(buf)?;
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct ListedConfigResource {
pub name: String,
pub resource_type: ConfigResourceType,
}
#[derive(Debug, Clone)]
pub struct ListConfigResourcesResponse {
pub throttle_time_ms: i32,
pub error_code: ErrorCode,
pub config_resources: Vec<ListedConfigResource>,
}
impl ListConfigResourcesResponse {
pub fn decode_v0(buf: &mut impl Buf) -> Result<Self> {
Self::decode(buf, false)
}
pub fn decode_v1(buf: &mut impl Buf) -> Result<Self> {
Self::decode(buf, true)
}
fn decode(buf: &mut impl Buf, with_resource_type: bool) -> Result<Self> {
let throttle_time_ms = i32::decode(buf)?;
let error_code = ErrorCode::from_i16(i16::decode(buf)?);
let resource_count =
check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut config_resources =
Vec::with_capacity(decode_capacity(resource_count, buf.remaining()));
for _ in 0..resource_count {
let name =
non_nullable_string("config resource name", KafkaString::decode_compact(buf)?.0)?;
let resource_type = if with_resource_type {
ConfigResourceType::from_i8(i8::decode(buf)?)
} else {
ConfigResourceType::ClientMetrics
};
let _ = TaggedFields::decode(buf)?;
config_resources.push(ListedConfigResource {
name,
resource_type,
});
}
let _ = TaggedFields::decode(buf)?;
Ok(Self {
throttle_time_ms,
error_code,
config_resources,
})
}
}
impl VersionedEncode for ListConfigResourcesRequest {
fn encode_versioned(&self, version: i16, buf: &mut impl BufMut) -> Result<()> {
match version {
0 => self.encode_v0(buf)?,
1 => self.encode_v1(buf)?,
_ => return unsupported_encode!("ListConfigResourcesRequest", version),
}
Ok(())
}
}
impl VersionedDecode for ListConfigResourcesResponse {
fn decode_versioned(version: i16, buf: &mut impl Buf) -> Result<Self> {
match version {
0 => Self::decode_v0(buf),
1 => Self::decode_v1(buf),
_ => unsupported_decode!("ListConfigResourcesResponse", version),
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::util::varint;
use bytes::BytesMut;
fn put_compact_string(buf: &mut BytesMut, s: Option<&str>) {
match s {
Some(val) => {
buf.put_u8((val.len() + 1) as u8);
buf.put_slice(val.as_bytes());
}
None => buf.put_u8(0),
}
}
fn put_tagged_fields(buf: &mut BytesMut) {
buf.put_u8(0);
}
#[test]
fn resource_type_wire_values_round_trip() {
for (ty, wire) in [
(ConfigResourceType::Group, 32i8),
(ConfigResourceType::ClientMetrics, 16),
(ConfigResourceType::BrokerLogger, 8),
(ConfigResourceType::Broker, 4),
(ConfigResourceType::Topic, 2),
] {
assert_eq!(ty.to_i8(), wire);
assert_eq!(ConfigResourceType::from_i8(wire), ty);
}
assert_eq!(ConfigResourceType::from_i8(64), ConfigResourceType::Unknown);
}
#[test]
fn request_encode_v0_is_tagged_fields_only() {
let req = ListConfigResourcesRequest::with_types(vec![ConfigResourceType::Topic]);
let mut buf = BytesMut::new();
req.encode_v0(&mut buf).unwrap();
let mut cur = &buf[..];
assert_eq!(cur.get_u8(), 0);
assert!(cur.is_empty());
}
#[test]
fn request_encode_v1_emits_resource_types() {
let req = ListConfigResourcesRequest::with_types(vec![
ConfigResourceType::Topic,
ConfigResourceType::Group,
]);
let mut buf = BytesMut::new();
req.encode_v1(&mut buf).unwrap();
let mut cur = &buf[..];
assert_eq!(varint::decode_unsigned_varint(&mut cur).unwrap(), 3); assert_eq!(cur.get_i8(), 2); assert_eq!(cur.get_i8(), 32); assert_eq!(cur.get_u8(), 0); assert!(cur.is_empty());
}
#[test]
fn request_encode_v1_empty_means_broker_default() {
let req = ListConfigResourcesRequest::all();
let mut buf = BytesMut::new();
req.encode_v1(&mut buf).unwrap();
let mut cur = &buf[..];
assert_eq!(varint::decode_unsigned_varint(&mut cur).unwrap(), 1); assert_eq!(cur.get_u8(), 0);
assert!(cur.is_empty());
}
#[test]
fn response_decode_v0_reports_client_metrics() {
let mut buf = BytesMut::new();
buf.put_i32(0); buf.put_i16(0); varint::encode_unsigned_varint(3, &mut buf); put_compact_string(&mut buf, Some("metric-a"));
put_tagged_fields(&mut buf);
put_compact_string(&mut buf, Some("metric-b"));
put_tagged_fields(&mut buf);
put_tagged_fields(&mut buf);
let resp = ListConfigResourcesResponse::decode_v0(&mut buf.freeze()).unwrap();
assert_eq!(resp.throttle_time_ms, 0);
assert!(resp.error_code.is_ok());
assert_eq!(resp.config_resources.len(), 2);
assert_eq!(resp.config_resources[0].name, "metric-a");
assert_eq!(
resp.config_resources[0].resource_type,
ConfigResourceType::ClientMetrics
);
assert_eq!(resp.config_resources[1].name, "metric-b");
}
#[test]
fn response_decode_v0_empty() {
let mut buf = BytesMut::new();
buf.put_i32(10); buf.put_i16(0); varint::encode_unsigned_varint(1, &mut buf); put_tagged_fields(&mut buf);
let resp = ListConfigResourcesResponse::decode_v0(&mut buf.freeze()).unwrap();
assert_eq!(resp.throttle_time_ms, 10);
assert!(resp.config_resources.is_empty());
}
#[test]
fn response_decode_v1_carries_resource_type() {
let mut buf = BytesMut::new();
buf.put_i32(7); buf.put_i16(0); varint::encode_unsigned_varint(4, &mut buf); put_compact_string(&mut buf, Some("my-topic"));
buf.put_i8(2); put_tagged_fields(&mut buf);
put_compact_string(&mut buf, Some("my-group"));
buf.put_i8(32); put_tagged_fields(&mut buf);
put_compact_string(&mut buf, Some("future"));
buf.put_i8(64); put_tagged_fields(&mut buf);
put_tagged_fields(&mut buf);
let resp = ListConfigResourcesResponse::decode_v1(&mut buf.freeze()).unwrap();
assert_eq!(resp.throttle_time_ms, 7);
assert_eq!(resp.config_resources.len(), 3);
assert_eq!(
resp.config_resources[0].resource_type,
ConfigResourceType::Topic
);
assert_eq!(
resp.config_resources[1].resource_type,
ConfigResourceType::Group
);
assert_eq!(
resp.config_resources[2].resource_type,
ConfigResourceType::Unknown
);
}
#[test]
fn versioned_dispatch_rejects_unknown_versions() {
let req = ListConfigResourcesRequest::all();
for v in [0, 1] {
let mut buf = BytesMut::new();
req.encode_versioned(v, &mut buf).unwrap();
assert!(!buf.is_empty());
}
let mut buf = BytesMut::new();
assert!(req.encode_versioned(2, &mut buf).is_err());
let mut empty = BytesMut::new().freeze();
assert!(ListConfigResourcesResponse::decode_versioned(2, &mut empty).is_err());
}
}