1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 90;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct DescribeShareGroupOffsetsTopic {
11 pub topic_name: String,
12 pub partitions: Vec<i32>,
13}
14
15#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct DescribeShareGroupOffsetsGroup {
18 pub group_id: String,
19 pub topics: Option<Vec<DescribeShareGroupOffsetsTopic>>,
20}
21
22fn encode_request(
23 correlation_id: i32,
24 client_id: Option<String>,
25 groups: &[DescribeShareGroupOffsetsGroup],
26 api_version: i16,
27) -> Result<Vec<u8>> {
28 let mut encoder = Encoder::new();
29 RequestHeader {
30 api_key: API_KEY,
31 api_version,
32 correlation_id,
33 client_id,
34 }
35 .encode_v2(&mut encoder)?;
36 encoder.write_compact_array(Some(groups), |encoder, group| {
37 encoder.write_compact_string(&group.group_id)?;
38 encoder.write_compact_array(group.topics.as_deref(), |encoder, topic| {
39 encoder.write_compact_string(&topic.topic_name)?;
40 encoder.write_compact_array(Some(&topic.partitions), |encoder, partition| {
41 encoder.write_i32(*partition);
42 Ok(())
43 })?;
44 encoder.write_empty_tagged_fields();
45 Ok(())
46 })?;
47 encoder.write_empty_tagged_fields();
48 Ok(())
49 })?;
50 encoder.write_empty_tagged_fields();
51 Ok(encoder.into_bytes())
52}
53
54#[derive(Debug, Clone, PartialEq, Eq)]
56pub struct DescribeShareGroupOffsetsRequestV0 {
57 pub correlation_id: i32,
58 pub client_id: Option<String>,
59 pub groups: Vec<DescribeShareGroupOffsetsGroup>,
60}
61
62impl DescribeShareGroupOffsetsRequestV0 {
63 pub fn encode(&self) -> Result<Vec<u8>> {
65 encode_request(self.correlation_id, self.client_id.clone(), &self.groups, 0)
66 }
67}
68
69#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct DescribeShareGroupOffsetsRequestV1 {
72 pub correlation_id: i32,
73 pub client_id: Option<String>,
74 pub groups: Vec<DescribeShareGroupOffsetsGroup>,
75}
76
77impl DescribeShareGroupOffsetsRequestV1 {
78 pub fn encode(&self) -> Result<Vec<u8>> {
80 encode_request(self.correlation_id, self.client_id.clone(), &self.groups, 1)
81 }
82}
83
84#[derive(Debug, Clone, PartialEq, Eq)]
86pub struct DescribeShareGroupOffsetsPartitionV0 {
87 pub partition_index: i32,
88 pub start_offset: i64,
89 pub leader_epoch: i32,
90 pub error_code: i16,
91 pub error_message: Option<String>,
92}
93
94#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct DescribeShareGroupOffsetsTopicResultV0 {
97 pub topic_name: String,
98 pub topic_id: [u8; 16],
99 pub partitions: Vec<DescribeShareGroupOffsetsPartitionV0>,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct DescribeShareGroupOffsetsGroupResultV0 {
105 pub group_id: String,
106 pub topics: Vec<DescribeShareGroupOffsetsTopicResultV0>,
107 pub error_code: i16,
108 pub error_message: Option<String>,
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
113pub struct DescribeShareGroupOffsetsResponseV0 {
114 pub throttle_time_ms: i32,
115 pub groups: Vec<DescribeShareGroupOffsetsGroupResultV0>,
116}
117
118impl DescribeShareGroupOffsetsResponseV0 {
119 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
121 let throttle_time_ms = decoder.read_i32()?;
122 let groups = decoder
123 .read_compact_array("share group offset groups", |decoder| {
124 decode_group_v0(decoder)
125 })?
126 .unwrap_or_default();
127 decoder.read_tagged_fields()?;
128 Ok(Self {
129 throttle_time_ms,
130 groups,
131 })
132 }
133}
134
135fn decode_group_v0(decoder: &mut Decoder<'_>) -> Result<DescribeShareGroupOffsetsGroupResultV0> {
136 let group_id = decoder.read_compact_string()?;
137 let topics = decoder
138 .read_compact_array("share group offset topics", |decoder| {
139 decode_topic_v0(decoder)
140 })?
141 .unwrap_or_default();
142 let error_code = decoder.read_i16()?;
143 let error_message = decoder.read_compact_nullable_string()?;
144 decoder.read_tagged_fields()?;
145 Ok(DescribeShareGroupOffsetsGroupResultV0 {
146 group_id,
147 topics,
148 error_code,
149 error_message,
150 })
151}
152
153fn decode_topic_v0(decoder: &mut Decoder<'_>) -> Result<DescribeShareGroupOffsetsTopicResultV0> {
154 let topic_name = decoder.read_compact_string()?;
155 let topic_id = decoder.read_uuid()?;
156 let partitions = decoder
157 .read_compact_array("share group offset partitions", |decoder| {
158 let partition = DescribeShareGroupOffsetsPartitionV0 {
159 partition_index: decoder.read_i32()?,
160 start_offset: decoder.read_i64()?,
161 leader_epoch: decoder.read_i32()?,
162 error_code: decoder.read_i16()?,
163 error_message: decoder.read_compact_nullable_string()?,
164 };
165 decoder.read_tagged_fields()?;
166 Ok(partition)
167 })?
168 .unwrap_or_default();
169 decoder.read_tagged_fields()?;
170 Ok(DescribeShareGroupOffsetsTopicResultV0 {
171 topic_name,
172 topic_id,
173 partitions,
174 })
175}
176
177#[derive(Debug, Clone, PartialEq, Eq)]
179pub struct DescribeShareGroupOffsetsPartitionV1 {
180 pub partition_index: i32,
181 pub start_offset: i64,
182 pub leader_epoch: i32,
183 pub lag: i64,
184 pub error_code: i16,
185 pub error_message: Option<String>,
186}
187
188#[derive(Debug, Clone, PartialEq, Eq)]
190pub struct DescribeShareGroupOffsetsTopicResultV1 {
191 pub topic_name: String,
192 pub topic_id: [u8; 16],
193 pub partitions: Vec<DescribeShareGroupOffsetsPartitionV1>,
194}
195
196#[derive(Debug, Clone, PartialEq, Eq)]
198pub struct DescribeShareGroupOffsetsGroupResultV1 {
199 pub group_id: String,
200 pub topics: Vec<DescribeShareGroupOffsetsTopicResultV1>,
201 pub error_code: i16,
202 pub error_message: Option<String>,
203}
204
205#[derive(Debug, Clone, PartialEq, Eq)]
207pub struct DescribeShareGroupOffsetsResponseV1 {
208 pub throttle_time_ms: i32,
209 pub groups: Vec<DescribeShareGroupOffsetsGroupResultV1>,
210}
211
212impl DescribeShareGroupOffsetsResponseV1 {
213 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
215 let throttle_time_ms = decoder.read_i32()?;
216 let groups = decoder
217 .read_compact_array("share group offset groups", |decoder| {
218 let group_id = decoder.read_compact_string()?;
219 let topics = decoder
220 .read_compact_array("share group offset topics", |decoder| {
221 let topic_name = decoder.read_compact_string()?;
222 let topic_id = decoder.read_uuid()?;
223 let partitions = decoder
224 .read_compact_array("share group offset partitions", |decoder| {
225 let result = DescribeShareGroupOffsetsPartitionV1 {
226 partition_index: decoder.read_i32()?,
227 start_offset: decoder.read_i64()?,
228 leader_epoch: decoder.read_i32()?,
229 lag: decoder.read_i64()?,
230 error_code: decoder.read_i16()?,
231 error_message: decoder.read_compact_nullable_string()?,
232 };
233 decoder.read_tagged_fields()?;
234 Ok(result)
235 })?
236 .unwrap_or_default();
237 decoder.read_tagged_fields()?;
238 Ok(DescribeShareGroupOffsetsTopicResultV1 {
239 topic_name,
240 topic_id,
241 partitions,
242 })
243 })?
244 .unwrap_or_default();
245 let error_code = decoder.read_i16()?;
246 let error_message = decoder.read_compact_nullable_string()?;
247 decoder.read_tagged_fields()?;
248 Ok(DescribeShareGroupOffsetsGroupResultV1 {
249 group_id,
250 topics,
251 error_code,
252 error_message,
253 })
254 })?
255 .unwrap_or_default();
256 decoder.read_tagged_fields()?;
257 Ok(Self {
258 throttle_time_ms,
259 groups,
260 })
261 }
262}
263
264#[cfg(test)]
265#[allow(clippy::unwrap_used)]
266mod tests {
267 use super::{
268 DescribeShareGroupOffsetsGroup, DescribeShareGroupOffsetsGroupResultV0,
269 DescribeShareGroupOffsetsPartitionV0, DescribeShareGroupOffsetsRequestV0,
270 DescribeShareGroupOffsetsResponseV0, DescribeShareGroupOffsetsTopic,
271 DescribeShareGroupOffsetsTopicResultV0, API_KEY,
272 };
273 use crate::codec::{Decoder, Encoder};
274
275 #[test]
276 fn encodes_describe_share_group_offsets_v0_request() {
277 let request = DescribeShareGroupOffsetsRequestV0 {
278 correlation_id: 23,
279 client_id: Some("kafrust".to_owned()),
280 groups: vec![DescribeShareGroupOffsetsGroup {
281 group_id: "share-orders".to_owned(),
282 topics: Some(vec![DescribeShareGroupOffsetsTopic {
283 topic_name: "orders".to_owned(),
284 partitions: vec![0, 2],
285 }]),
286 }],
287 };
288
289 let encoded = request.encode().unwrap();
290 assert_eq!(&encoded[..4], &[0, 90, 0, 0]);
291 assert_eq!(API_KEY, 90);
292 assert_eq!(encoded.last(), Some(&0));
293 }
294
295 #[test]
296 fn decodes_describe_share_group_offsets_v0_response() -> crate::error::Result<()> {
297 let mut bytes = Encoder::new();
298 bytes.write_i32(12);
299 bytes.write_compact_array(Some(&[()]), |encoder, ()| {
300 encoder.write_compact_string("share-orders")?;
301 encoder.write_compact_array(Some(&[()]), |encoder, ()| {
302 encoder.write_compact_string("orders")?;
303 encoder.write_uuid(&[7; 16]);
304 encoder.write_compact_array(Some(&[()]), |encoder, ()| {
305 encoder.write_i32(0);
306 encoder.write_i64(42);
307 encoder.write_i32(3);
308 encoder.write_i16(0);
309 encoder.write_compact_nullable_string(None)?;
310 encoder.write_empty_tagged_fields();
311 Ok(())
312 })?;
313 encoder.write_empty_tagged_fields();
314 Ok(())
315 })?;
316 encoder.write_i16(0);
317 encoder.write_compact_nullable_string(None)?;
318 encoder.write_empty_tagged_fields();
319 Ok(())
320 })?;
321 bytes.write_empty_tagged_fields();
322
323 let encoded = bytes.into_bytes();
324 let mut decoder = Decoder::new(&encoded);
325 let response = DescribeShareGroupOffsetsResponseV0::decode_body(&mut decoder)?;
326 assert_eq!(response.throttle_time_ms, 12);
327 assert_eq!(
328 response.groups,
329 vec![DescribeShareGroupOffsetsGroupResultV0 {
330 group_id: "share-orders".to_owned(),
331 topics: vec![DescribeShareGroupOffsetsTopicResultV0 {
332 topic_name: "orders".to_owned(),
333 topic_id: [7; 16],
334 partitions: vec![DescribeShareGroupOffsetsPartitionV0 {
335 partition_index: 0,
336 start_offset: 42,
337 leader_epoch: 3,
338 error_code: 0,
339 error_message: None,
340 }],
341 }],
342 error_code: 0,
343 error_message: None,
344 }]
345 );
346 assert!(decoder.is_empty());
347 Ok(())
348 }
349
350 #[test]
351 fn decodes_describe_share_group_offsets_v1_lag() -> crate::error::Result<()> {
352 let mut bytes = Encoder::new();
353 bytes.write_i32(0);
354 bytes.write_compact_array(Some(&[()]), |encoder, ()| {
355 encoder.write_compact_string("share-orders")?;
356 encoder.write_compact_array(Some(&[()]), |encoder, ()| {
357 encoder.write_compact_string("orders")?;
358 encoder.write_uuid(&[9; 16]);
359 encoder.write_compact_array(Some(&[()]), |encoder, ()| {
360 encoder.write_i32(1);
361 encoder.write_i64(100);
362 encoder.write_i32(4);
363 encoder.write_i64(7);
364 encoder.write_i16(0);
365 encoder.write_compact_nullable_string(None)?;
366 encoder.write_empty_tagged_fields();
367 Ok(())
368 })?;
369 encoder.write_empty_tagged_fields();
370 Ok(())
371 })?;
372 encoder.write_i16(0);
373 encoder.write_compact_nullable_string(None)?;
374 encoder.write_empty_tagged_fields();
375 Ok(())
376 })?;
377 bytes.write_empty_tagged_fields();
378
379 let encoded = bytes.into_bytes();
380 let mut decoder = Decoder::new(&encoded);
381 let response = super::DescribeShareGroupOffsetsResponseV1::decode_body(&mut decoder)?;
382 assert_eq!(response.groups[0].topics[0].partitions[0].lag, 7);
383 assert!(decoder.is_empty());
384 Ok(())
385 }
386}