1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 60;
7
8pub const BROKER_ENDPOINT_TYPE: i8 = 1;
10pub const CONTROLLER_ENDPOINT_TYPE: i8 = 2;
12
13#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct DescribeClusterRequest {
16 pub correlation_id: i32,
17 pub client_id: Option<String>,
18 pub api_version: i16,
19 pub include_cluster_authorized_operations: bool,
20 pub endpoint_type: i8,
21}
22
23impl DescribeClusterRequest {
24 pub fn encode(&self) -> Result<Vec<u8>> {
26 let mut encoder = Encoder::new();
27 RequestHeader {
28 api_key: API_KEY,
29 api_version: self.api_version,
30 correlation_id: self.correlation_id,
31 client_id: self.client_id.clone(),
32 }
33 .encode_v2(&mut encoder)?;
34 encoder.write_bool(self.include_cluster_authorized_operations);
35 if self.api_version >= 1 {
36 encoder.write_i8(self.endpoint_type);
37 }
38 encoder.write_empty_tagged_fields();
39 Ok(encoder.into_bytes())
40 }
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
45pub struct DescribeClusterBroker {
46 pub node_id: i32,
47 pub host: String,
48 pub port: i32,
49 pub rack: Option<String>,
50}
51
52#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct DescribeClusterResponse {
55 pub api_version: i16,
56 pub throttle_time_ms: i32,
57 pub error_code: i16,
58 pub error_message: Option<String>,
59 pub endpoint_type: Option<i8>,
60 pub cluster_id: String,
61 pub controller_id: i32,
62 pub brokers: Vec<DescribeClusterBroker>,
63 pub cluster_authorized_operations: i32,
64}
65
66impl DescribeClusterResponse {
67 pub fn decode_body(decoder: &mut Decoder<'_>, api_version: i16) -> Result<Self> {
69 let throttle_time_ms = decoder.read_i32()?;
70 let error_code = decoder.read_i16()?;
71 let error_message = decoder.read_compact_nullable_string()?;
72 let endpoint_type = (api_version >= 1).then(|| decoder.read_i8()).transpose()?;
73 let cluster_id = decoder.read_compact_string()?;
74 let controller_id = decoder.read_i32()?;
75 let brokers = decoder
76 .read_compact_array("describe cluster brokers", |decoder| {
77 let node_id = decoder.read_i32()?;
78 let host = decoder.read_compact_string()?;
79 let port = decoder.read_i32()?;
80 let rack = decoder.read_compact_nullable_string()?;
81 decoder.read_tagged_fields()?;
82 Ok(DescribeClusterBroker {
83 node_id,
84 host,
85 port,
86 rack,
87 })
88 })?
89 .unwrap_or_default();
90 let cluster_authorized_operations = decoder.read_i32()?;
91 decoder.read_tagged_fields()?;
92 Ok(Self {
93 api_version,
94 throttle_time_ms,
95 error_code,
96 error_message,
97 endpoint_type,
98 cluster_id,
99 controller_id,
100 brokers,
101 cluster_authorized_operations,
102 })
103 }
104}
105
106#[cfg(test)]
107#[allow(clippy::unwrap_used)]
108mod tests {
109 use super::{DescribeClusterRequest, DescribeClusterResponse, API_KEY};
110 use crate::codec::{Decoder, Encoder};
111
112 #[test]
113 fn encodes_describe_cluster_v1_request() {
114 let request = DescribeClusterRequest {
115 correlation_id: 7,
116 client_id: Some("kafrust".to_owned()),
117 api_version: 1,
118 include_cluster_authorized_operations: true,
119 endpoint_type: 2,
120 };
121
122 assert_eq!(
123 request.encode().unwrap(),
124 [
125 0, 60, 0, 1, 0, 0, 0, 7, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',
129 0, 1, 2, 0, ]
134 );
135 assert_eq!(API_KEY, 60);
136 }
137
138 #[test]
139 fn decodes_describe_cluster_v1_response() {
140 let mut body = Encoder::new();
141 body.write_i32(11);
142 body.write_i16(0);
143 body.write_compact_nullable_string(None).unwrap();
144 body.write_i8(2);
145 body.write_compact_string("cluster").unwrap();
146 body.write_i32(2);
147 body.write_compact_array(Some(&["broker-a".to_owned()]), |encoder, host| {
148 encoder.write_i32(1);
149 encoder.write_compact_string(host)?;
150 encoder.write_i32(9092);
151 encoder.write_compact_nullable_string(Some("rack-a"))?;
152 encoder.write_empty_tagged_fields();
153 Ok(())
154 })
155 .unwrap();
156 body.write_i32(0x1234);
157 body.write_empty_tagged_fields();
158
159 let bytes = body.into_bytes();
160 let mut decoder = Decoder::new(&bytes);
161 let response = DescribeClusterResponse::decode_body(&mut decoder, 1).unwrap();
162
163 assert_eq!(response.api_version, 1);
164 assert_eq!(response.throttle_time_ms, 11);
165 assert_eq!(response.endpoint_type, Some(2));
166 assert_eq!(response.cluster_id, "cluster");
167 assert_eq!(response.controller_id, 2);
168 assert_eq!(response.brokers.len(), 1);
169 assert_eq!(response.brokers[0].host, "broker-a");
170 assert_eq!(response.brokers[0].rack.as_deref(), Some("rack-a"));
171 assert_eq!(response.cluster_authorized_operations, 0x1234);
172 assert!(decoder.is_empty());
173 }
174}