Skip to main content

kafrust_protocol/api/
consumer_group_describe.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka ConsumerGroupDescribe API key.
6pub 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    /// -1 is unknown, 0 is classic, and 1 is a consumer-protocol member.
171    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}