1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 77;
7
8#[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 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#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct ShareGroupDescribeResponseV1 {
43 pub throttle_time_ms: i32,
44 pub groups: Vec<DescribedShareGroup>,
45}
46
47impl ShareGroupDescribeResponseV1 {
48 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#[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#[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#[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#[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}