Skip to main content

kafrust_protocol/api/
describe_share_group_offsets.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka DescribeShareGroupOffsets API key.
6pub const API_KEY: i16 = 90;
7
8/// One topic and partition filter for a share-group offset query.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct DescribeShareGroupOffsetsTopic {
11    pub topic_name: String,
12    pub partitions: Vec<i32>,
13}
14
15/// One share group in a DescribeShareGroupOffsets request.
16#[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/// DescribeShareGroupOffsets v0 request.
55#[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    /// Encodes the flexible request, including its request header.
64    pub fn encode(&self) -> Result<Vec<u8>> {
65        encode_request(self.correlation_id, self.client_id.clone(), &self.groups, 0)
66    }
67}
68
69/// DescribeShareGroupOffsets v1 request.
70#[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    /// Encodes the flexible request, including its request header.
79    pub fn encode(&self) -> Result<Vec<u8>> {
80        encode_request(self.correlation_id, self.client_id.clone(), &self.groups, 1)
81    }
82}
83
84/// One partition result returned by DescribeShareGroupOffsets v0.
85#[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/// One topic result returned by DescribeShareGroupOffsets v0.
95#[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/// One group result returned by DescribeShareGroupOffsets v0.
103#[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/// DescribeShareGroupOffsets v0 response.
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub struct DescribeShareGroupOffsetsResponseV0 {
114    pub throttle_time_ms: i32,
115    pub groups: Vec<DescribeShareGroupOffsetsGroupResultV0>,
116}
117
118impl DescribeShareGroupOffsetsResponseV0 {
119    /// Decodes the flexible response body after the response header.
120    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/// One partition result returned by DescribeShareGroupOffsets v1.
178#[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/// One topic result returned by DescribeShareGroupOffsets v1.
189#[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/// One group result returned by DescribeShareGroupOffsets v1.
197#[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/// DescribeShareGroupOffsets v1 response.
206#[derive(Debug, Clone, PartialEq, Eq)]
207pub struct DescribeShareGroupOffsetsResponseV1 {
208    pub throttle_time_ms: i32,
209    pub groups: Vec<DescribeShareGroupOffsetsGroupResultV1>,
210}
211
212impl DescribeShareGroupOffsetsResponseV1 {
213    /// Decodes the flexible response body after the response header.
214    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}