Skip to main content

kafrust_protocol/api/
share_group_describe.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka ShareGroupDescribe API key.
6pub const API_KEY: i16 = 77;
7
8/// ShareGroupDescribe v1 request.
9///
10/// Version 1 is the stable KIP-932 wire shape. Kafka 4.0's early-access v0
11/// was removed in Kafka 4.1, so the public protocol surface starts at v1.
12#[derive(Debug, Clone, PartialEq, Eq)]
13pub struct ShareGroupDescribeRequestV1 {
14    pub correlation_id: i32,
15    pub client_id: Option<String>,
16    pub group_ids: Vec<String>,
17    pub include_authorized_operations: bool,
18}
19
20impl ShareGroupDescribeRequestV1 {
21    /// Encodes the flexible request, including its request header.
22    pub fn encode(&self) -> Result<Vec<u8>> {
23        let mut encoder = Encoder::new();
24        RequestHeader {
25            api_key: API_KEY,
26            api_version: 1,
27            correlation_id: self.correlation_id,
28            client_id: self.client_id.clone(),
29        }
30        .encode_v2(&mut encoder)?;
31        encoder.write_compact_array(Some(&self.group_ids), |encoder, group_id| {
32            encoder.write_compact_string(group_id)
33        })?;
34        encoder.write_bool(self.include_authorized_operations);
35        encoder.write_empty_tagged_fields();
36        Ok(encoder.into_bytes())
37    }
38}
39
40/// ShareGroupDescribe v1 response.
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct ShareGroupDescribeResponseV1 {
43    pub throttle_time_ms: i32,
44    pub groups: Vec<DescribedShareGroup>,
45}
46
47impl ShareGroupDescribeResponseV1 {
48    /// Decodes the flexible response body after the response header.
49    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
50        let throttle_time_ms = decoder.read_i32()?;
51        let groups = decoder
52            .read_compact_array("share group descriptions", DescribedShareGroup::decode)?
53            .unwrap_or_default();
54        decoder.read_tagged_fields()?;
55        Ok(Self {
56            throttle_time_ms,
57            groups,
58        })
59    }
60}
61
62/// One share group returned by ShareGroupDescribe.
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct DescribedShareGroup {
65    pub error_code: i16,
66    pub error_message: Option<String>,
67    pub group_id: String,
68    pub group_state: String,
69    pub group_epoch: i32,
70    pub assignment_epoch: i32,
71    pub assignor_name: String,
72    pub members: Vec<DescribedShareGroupMember>,
73    pub authorized_operations: i32,
74}
75
76impl DescribedShareGroup {
77    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
78        let error_code = decoder.read_i16()?;
79        let error_message = decoder.read_compact_nullable_string()?;
80        let group_id = decoder.read_compact_string()?;
81        let group_state = decoder.read_compact_string()?;
82        let group_epoch = decoder.read_i32()?;
83        let assignment_epoch = decoder.read_i32()?;
84        let assignor_name = decoder.read_compact_string()?;
85        let members = decoder
86            .read_compact_array("share group members", DescribedShareGroupMember::decode)?
87            .unwrap_or_default();
88        let authorized_operations = decoder.read_i32()?;
89        decoder.read_tagged_fields()?;
90        Ok(Self {
91            error_code,
92            error_message,
93            group_id,
94            group_state,
95            group_epoch,
96            assignment_epoch,
97            assignor_name,
98            members,
99            authorized_operations,
100        })
101    }
102}
103
104/// One member returned by ShareGroupDescribe.
105#[derive(Debug, Clone, PartialEq, Eq)]
106pub struct DescribedShareGroupMember {
107    pub member_id: String,
108    pub rack_id: Option<String>,
109    pub member_epoch: i32,
110    pub client_id: String,
111    pub client_host: String,
112    pub subscribed_topic_names: Vec<String>,
113    pub assignment: ShareGroupDescribeAssignment,
114}
115
116impl DescribedShareGroupMember {
117    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
118        let member_id = decoder.read_compact_string()?;
119        let rack_id = decoder.read_compact_nullable_string()?;
120        let member_epoch = decoder.read_i32()?;
121        let client_id = decoder.read_compact_string()?;
122        let client_host = decoder.read_compact_string()?;
123        let subscribed_topic_names = decoder
124            .read_compact_array("share group subscribed topics", |decoder| {
125                decoder.read_compact_string()
126            })?
127            .unwrap_or_default();
128        let assignment = ShareGroupDescribeAssignment::decode(decoder)?;
129        decoder.read_tagged_fields()?;
130        Ok(Self {
131            member_id,
132            rack_id,
133            member_epoch,
134            client_id,
135            client_host,
136            subscribed_topic_names,
137            assignment,
138        })
139    }
140}
141
142/// Assignment returned for one share-group member.
143#[derive(Debug, Clone, PartialEq, Eq)]
144pub struct ShareGroupDescribeAssignment {
145    pub topic_partitions: Vec<ShareGroupDescribeTopicPartitions>,
146}
147
148impl ShareGroupDescribeAssignment {
149    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
150        let topic_partitions = decoder
151            .read_compact_array("share group assignment topics", |decoder| {
152                ShareGroupDescribeTopicPartitions::decode(decoder)
153            })?
154            .unwrap_or_default();
155        decoder.read_tagged_fields()?;
156        Ok(Self { topic_partitions })
157    }
158}
159
160/// Topic partitions assigned to one share-group member.
161#[derive(Debug, Clone, PartialEq, Eq)]
162pub struct ShareGroupDescribeTopicPartitions {
163    pub topic_id: [u8; 16],
164    pub topic_name: String,
165    pub partitions: Vec<i32>,
166}
167
168impl ShareGroupDescribeTopicPartitions {
169    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
170        let topic_id = decoder.read_uuid()?;
171        let topic_name = decoder.read_compact_string()?;
172        let partitions = decoder
173            .read_compact_array("share group assignment partitions", |decoder| {
174                decoder.read_i32()
175            })?
176            .unwrap_or_default();
177        decoder.read_tagged_fields()?;
178        Ok(Self {
179            topic_id,
180            topic_name,
181            partitions,
182        })
183    }
184}
185
186#[cfg(test)]
187#[allow(clippy::unwrap_used)]
188mod tests {
189    use super::{ShareGroupDescribeRequestV1, ShareGroupDescribeResponseV1, API_KEY};
190    use crate::codec::{Decoder, Encoder};
191
192    #[test]
193    fn encodes_share_group_describe_v1_request() {
194        let request = ShareGroupDescribeRequestV1 {
195            correlation_id: 23,
196            client_id: Some("kafrust".to_owned()),
197            group_ids: vec!["share-orders".to_owned()],
198            include_authorized_operations: true,
199        };
200
201        let encoded = request.encode().unwrap();
202        assert_eq!(&encoded[..4], &[0, 77, 0, 1]);
203        assert_eq!(API_KEY, 77);
204    }
205
206    #[test]
207    fn decodes_share_group_describe_v1_response() -> crate::error::Result<()> {
208        let mut bytes = Encoder::new();
209        bytes.write_i32(12);
210        bytes.write_compact_array(Some(&[1_i8]), |encoder, _| {
211            encoder.write_i16(0);
212            encoder.write_compact_nullable_string(Some("ok"))?;
213            encoder.write_compact_string("share-orders")?;
214            encoder.write_compact_string("Stable")?;
215            encoder.write_i32(4);
216            encoder.write_i32(5);
217            encoder.write_compact_string("uniform")?;
218            encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
219                encoder.write_compact_string("member-1")?;
220                encoder.write_compact_nullable_string(Some("rack-a"))?;
221                encoder.write_i32(7);
222                encoder.write_compact_string("client-a")?;
223                encoder.write_compact_string("/127.0.0.1")?;
224                encoder.write_compact_array(Some(&["orders".to_owned()]), |encoder, topic| {
225                    encoder.write_compact_string(topic)
226                })?;
227                encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
228                    encoder.write_uuid(&[7; 16]);
229                    encoder.write_compact_string("orders")?;
230                    encoder.write_compact_array(Some(&[0_i32, 2]), |encoder, partition| {
231                        encoder.write_i32(*partition);
232                        Ok(())
233                    })?;
234                    encoder.write_empty_tagged_fields();
235                    Ok(())
236                })?;
237                encoder.write_empty_tagged_fields();
238                encoder.write_empty_tagged_fields();
239                Ok(())
240            })?;
241            encoder.write_i32(-2147483648);
242            encoder.write_empty_tagged_fields();
243            Ok(())
244        })?;
245        bytes.write_empty_tagged_fields();
246
247        let encoded = bytes.into_bytes();
248        let mut decoder = Decoder::new(&encoded);
249        let response = ShareGroupDescribeResponseV1::decode_body(&mut decoder)?;
250
251        assert_eq!(response.throttle_time_ms, 12);
252        assert_eq!(response.groups[0].group_id, "share-orders");
253        assert_eq!(response.groups[0].group_epoch, 4);
254        assert_eq!(response.groups[0].members[0].member_epoch, 7);
255        assert_eq!(
256            response.groups[0].members[0].assignment.topic_partitions[0].partitions,
257            [0, 2]
258        );
259        assert_eq!(response.groups[0].authorized_operations, -2147483648);
260        assert!(decoder.is_empty());
261        Ok(())
262    }
263}