1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 69;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct ConsumerGroupDescribeRequestV0 {
10 pub correlation_id: i32,
11 pub client_id: Option<String>,
12 pub group_ids: Vec<String>,
13 pub include_authorized_operations: bool,
14}
15
16impl ConsumerGroupDescribeRequestV0 {
17 pub fn encode(&self) -> Result<Vec<u8>> {
18 encode_request(
19 0,
20 self.correlation_id,
21 self.client_id.as_deref(),
22 &self.group_ids,
23 self.include_authorized_operations,
24 )
25 }
26}
27
28#[derive(Debug, Clone, PartialEq, Eq)]
29pub struct ConsumerGroupDescribeRequestV1 {
30 pub correlation_id: i32,
31 pub client_id: Option<String>,
32 pub group_ids: Vec<String>,
33 pub include_authorized_operations: bool,
34}
35
36impl ConsumerGroupDescribeRequestV1 {
37 pub fn encode(&self) -> Result<Vec<u8>> {
38 encode_request(
39 1,
40 self.correlation_id,
41 self.client_id.as_deref(),
42 &self.group_ids,
43 self.include_authorized_operations,
44 )
45 }
46}
47
48fn encode_request(
49 api_version: i16,
50 correlation_id: i32,
51 client_id: Option<&str>,
52 group_ids: &[String],
53 include_authorized_operations: bool,
54) -> Result<Vec<u8>> {
55 let mut encoder = Encoder::new();
56 RequestHeader {
57 api_key: API_KEY,
58 api_version,
59 correlation_id,
60 client_id: client_id.map(str::to_owned),
61 }
62 .encode_v2(&mut encoder)?;
63 encoder.write_compact_array(Some(group_ids), |encoder, group_id| {
64 encoder.write_compact_string(group_id)
65 })?;
66 encoder.write_bool(include_authorized_operations);
67 encoder.write_empty_tagged_fields();
68 Ok(encoder.into_bytes())
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct ConsumerGroupDescribeResponseV0 {
73 pub throttle_time_ms: i32,
74 pub groups: Vec<DescribedConsumerGroup>,
75}
76
77impl ConsumerGroupDescribeResponseV0 {
78 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
79 let throttle_time_ms = decoder.read_i32()?;
80 let groups = decoder
81 .read_compact_array("consumer group descriptions", |decoder| {
82 DescribedConsumerGroup::decode(decoder, false)
83 })?
84 .unwrap_or_default();
85 decoder.read_tagged_fields()?;
86 Ok(Self {
87 throttle_time_ms,
88 groups,
89 })
90 }
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct ConsumerGroupDescribeResponseV1 {
95 pub throttle_time_ms: i32,
96 pub groups: Vec<DescribedConsumerGroup>,
97}
98
99impl ConsumerGroupDescribeResponseV1 {
100 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
101 let throttle_time_ms = decoder.read_i32()?;
102 let groups = decoder
103 .read_compact_array("consumer group descriptions", |decoder| {
104 DescribedConsumerGroup::decode(decoder, true)
105 })?
106 .unwrap_or_default();
107 decoder.read_tagged_fields()?;
108 Ok(Self {
109 throttle_time_ms,
110 groups,
111 })
112 }
113}
114
115#[derive(Debug, Clone, PartialEq, Eq)]
116pub struct DescribedConsumerGroup {
117 pub error_code: i16,
118 pub error_message: Option<String>,
119 pub group_id: String,
120 pub group_state: String,
121 pub group_epoch: i32,
122 pub assignment_epoch: i32,
123 pub assignor_name: String,
124 pub members: Vec<DescribedConsumerGroupMember>,
125 pub authorized_operations: i32,
126}
127
128impl DescribedConsumerGroup {
129 fn decode(decoder: &mut Decoder<'_>, has_member_type: bool) -> Result<Self> {
130 let error_code = decoder.read_i16()?;
131 let error_message = decoder.read_compact_nullable_string()?;
132 let group_id = decoder.read_compact_string()?;
133 let group_state = decoder.read_compact_string()?;
134 let group_epoch = decoder.read_i32()?;
135 let assignment_epoch = decoder.read_i32()?;
136 let assignor_name = decoder.read_compact_string()?;
137 let members = decoder
138 .read_compact_array("consumer group members", |decoder| {
139 DescribedConsumerGroupMember::decode(decoder, has_member_type)
140 })?
141 .unwrap_or_default();
142 let authorized_operations = decoder.read_i32()?;
143 decoder.read_tagged_fields()?;
144 Ok(Self {
145 error_code,
146 error_message,
147 group_id,
148 group_state,
149 group_epoch,
150 assignment_epoch,
151 assignor_name,
152 members,
153 authorized_operations,
154 })
155 }
156}
157
158#[derive(Debug, Clone, PartialEq, Eq)]
159pub struct DescribedConsumerGroupMember {
160 pub member_id: String,
161 pub instance_id: Option<String>,
162 pub rack_id: Option<String>,
163 pub member_epoch: i32,
164 pub client_id: String,
165 pub client_host: String,
166 pub subscribed_topic_names: Vec<String>,
167 pub subscribed_topic_regex: Option<String>,
168 pub assignment: ConsumerGroupDescribeAssignment,
169 pub target_assignment: ConsumerGroupDescribeAssignment,
170 pub member_type: i8,
172}
173
174impl DescribedConsumerGroupMember {
175 fn decode(decoder: &mut Decoder<'_>, has_member_type: bool) -> Result<Self> {
176 let member_id = decoder.read_compact_string()?;
177 let instance_id = decoder.read_compact_nullable_string()?;
178 let rack_id = decoder.read_compact_nullable_string()?;
179 let member_epoch = decoder.read_i32()?;
180 let client_id = decoder.read_compact_string()?;
181 let client_host = decoder.read_compact_string()?;
182 let subscribed_topic_names = decoder
183 .read_compact_array("subscribed topic names", |decoder| {
184 decoder.read_compact_string()
185 })?
186 .unwrap_or_default();
187 let subscribed_topic_regex = decoder.read_compact_nullable_string()?;
188 let assignment = ConsumerGroupDescribeAssignment::decode(decoder)?;
189 let target_assignment = ConsumerGroupDescribeAssignment::decode(decoder)?;
190 let member_type = if has_member_type {
191 decoder.read_i8()?
192 } else {
193 -1
194 };
195 decoder.read_tagged_fields()?;
196 Ok(Self {
197 member_id,
198 instance_id,
199 rack_id,
200 member_epoch,
201 client_id,
202 client_host,
203 subscribed_topic_names,
204 subscribed_topic_regex,
205 assignment,
206 target_assignment,
207 member_type,
208 })
209 }
210}
211
212#[derive(Debug, Clone, PartialEq, Eq)]
213pub struct ConsumerGroupDescribeAssignment {
214 pub topic_partitions: Vec<ConsumerGroupDescribeTopicPartitions>,
215}
216
217impl ConsumerGroupDescribeAssignment {
218 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
219 let topic_partitions = decoder
220 .read_compact_array("consumer group assignment topics", |decoder| {
221 ConsumerGroupDescribeTopicPartitions::decode(decoder)
222 })?
223 .unwrap_or_default();
224 decoder.read_tagged_fields()?;
225 Ok(Self { topic_partitions })
226 }
227}
228
229#[derive(Debug, Clone, PartialEq, Eq)]
230pub struct ConsumerGroupDescribeTopicPartitions {
231 pub topic_id: [u8; 16],
232 pub topic_name: String,
233 pub partitions: Vec<i32>,
234}
235
236impl ConsumerGroupDescribeTopicPartitions {
237 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
238 let topic_id = decoder.read_uuid()?;
239 let topic_name = decoder.read_compact_string()?;
240 let partitions = decoder
241 .read_compact_array("consumer group assignment partitions", |decoder| {
242 decoder.read_i32()
243 })?
244 .unwrap_or_default();
245 decoder.read_tagged_fields()?;
246 Ok(Self {
247 topic_id,
248 topic_name,
249 partitions,
250 })
251 }
252}
253
254#[cfg(test)]
255#[allow(clippy::unwrap_used)]
256mod tests {
257 use super::{ConsumerGroupDescribeRequestV1, ConsumerGroupDescribeResponseV1, API_KEY};
258 use crate::codec::{Decoder, Encoder};
259
260 #[test]
261 fn encodes_consumer_group_describe_v1_request() {
262 let request = ConsumerGroupDescribeRequestV1 {
263 correlation_id: 23,
264 client_id: Some("kafrust".to_owned()),
265 group_ids: vec!["orders".to_owned(), "payments".to_owned()],
266 include_authorized_operations: true,
267 };
268
269 let encoded = request.encode().unwrap();
270 assert_eq!(&encoded[..4], &[0, 69, 0, 1]);
271 assert_eq!(API_KEY, 69);
272 }
273
274 #[test]
275 fn decodes_consumer_group_describe_v1_response() -> crate::error::Result<()> {
276 let mut bytes = Encoder::new();
277 bytes.write_i32(9);
278 bytes.write_compact_array(Some(&[1_i8]), |encoder, _| {
279 encoder.write_i16(0);
280 encoder.write_compact_nullable_string(Some("ok"))?;
281 encoder.write_compact_string("orders")?;
282 encoder.write_compact_string("Stable")?;
283 encoder.write_i32(4);
284 encoder.write_i32(5);
285 encoder.write_compact_string("uniform")?;
286 encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
287 encoder.write_compact_string("member-1")?;
288 encoder.write_compact_nullable_string(None)?;
289 encoder.write_compact_nullable_string(Some("rack-a"))?;
290 encoder.write_i32(7);
291 encoder.write_compact_string("client-a")?;
292 encoder.write_compact_string("/127.0.0.1")?;
293 encoder.write_compact_array(Some(&["orders".to_owned()]), |encoder, topic| {
294 encoder.write_compact_string(topic)
295 })?;
296 encoder.write_compact_nullable_string(None)?;
297 for _ in 0..2 {
298 encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
299 encoder.write_uuid(&[7; 16]);
300 encoder.write_compact_string("orders")?;
301 encoder.write_compact_array(Some(&[0_i32, 2]), |encoder, partition| {
302 encoder.write_i32(*partition);
303 Ok(())
304 })?;
305 encoder.write_empty_tagged_fields();
306 Ok(())
307 })?;
308 encoder.write_empty_tagged_fields();
309 }
310 encoder.write_i8(1);
311 encoder.write_empty_tagged_fields();
312 Ok(())
313 })?;
314 encoder.write_i32(-2147483648);
315 encoder.write_empty_tagged_fields();
316 Ok(())
317 })?;
318 bytes.write_empty_tagged_fields();
319 let bytes = bytes.into_bytes();
320 let mut decoder = Decoder::new(&bytes);
321 let response = ConsumerGroupDescribeResponseV1::decode_body(&mut decoder).unwrap();
322
323 assert_eq!(response.throttle_time_ms, 9);
324 assert_eq!(response.groups[0].group_id, "orders");
325 assert_eq!(response.groups[0].group_epoch, 4);
326 assert_eq!(response.groups[0].members[0].member_type, 1);
327 assert_eq!(
328 response.groups[0].members[0].assignment.topic_partitions[0].partitions,
329 [0, 2]
330 );
331 assert_eq!(response.groups[0].authorized_operations, -2147483648);
332 assert!(decoder.is_empty());
333 Ok(())
334 }
335}