Skip to main content

kafrust_protocol/api/
offset_fetch.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 9;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct OffsetFetchRequestV2 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub group_id: String,
12    pub topics: Option<Vec<OffsetFetchTopic>>,
13}
14
15impl OffsetFetchRequestV2 {
16    pub fn encode(&self) -> Result<Vec<u8>> {
17        let mut encoder = Encoder::new();
18        RequestHeader {
19            api_key: API_KEY,
20            api_version: 2,
21            correlation_id: self.correlation_id,
22            client_id: self.client_id.clone(),
23        }
24        .encode_v1(&mut encoder)?;
25        encoder.write_string(&self.group_id)?;
26        encoder.write_array(self.topics.as_deref(), |encoder, topic| {
27            topic.encode(encoder)
28        })?;
29        Ok(encoder.into_bytes())
30    }
31}
32
33/// OffsetFetch v9 request for Kafka's KIP-848 consumer group protocol.
34#[derive(Debug, Clone, PartialEq, Eq)]
35pub struct OffsetFetchRequestV9 {
36    pub correlation_id: i32,
37    pub client_id: Option<String>,
38    pub group_id: String,
39    pub member_id: Option<String>,
40    pub member_epoch: i32,
41    pub topics: Option<Vec<OffsetFetchTopicV9>>,
42    pub require_stable: bool,
43}
44
45impl OffsetFetchRequestV9 {
46    pub fn encode(&self) -> Result<Vec<u8>> {
47        let mut encoder = Encoder::new();
48        RequestHeader {
49            api_key: API_KEY,
50            api_version: 9,
51            correlation_id: self.correlation_id,
52            client_id: self.client_id.clone(),
53        }
54        .encode_v2(&mut encoder)?;
55        encoder.write_compact_array(Some(&[()]), |encoder, ()| {
56            encoder.write_compact_string(&self.group_id)?;
57            encoder.write_compact_nullable_string(self.member_id.as_deref())?;
58            encoder.write_i32(self.member_epoch);
59            encoder.write_compact_array(self.topics.as_deref(), |encoder, topic| {
60                topic.encode(encoder)
61            })?;
62            encoder.write_empty_tagged_fields();
63            Ok(())
64        })?;
65        encoder.write_bool(self.require_stable);
66        encoder.write_empty_tagged_fields();
67        Ok(encoder.into_bytes())
68    }
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct OffsetFetchTopic {
73    pub name: String,
74    pub partition_indexes: Vec<i32>,
75}
76
77impl OffsetFetchTopic {
78    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
79        encoder.write_string(&self.name)?;
80        encoder.write_array(
81            Some(self.partition_indexes.as_slice()),
82            |encoder, partition| {
83                encoder.write_i32(*partition);
84                Ok(())
85            },
86        )
87    }
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
91pub struct OffsetFetchTopicV9 {
92    pub name: String,
93    pub partition_indexes: Vec<i32>,
94}
95
96impl OffsetFetchTopicV9 {
97    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
98        encoder.write_compact_string(&self.name)?;
99        encoder.write_compact_array(
100            Some(self.partition_indexes.as_slice()),
101            |encoder, partition| {
102                encoder.write_i32(*partition);
103                Ok(())
104            },
105        )?;
106        encoder.write_empty_tagged_fields();
107        Ok(())
108    }
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
112pub struct OffsetFetchResponseV2 {
113    pub topics: Vec<OffsetFetchTopicResponse>,
114    pub error_code: i16,
115}
116
117impl OffsetFetchResponseV2 {
118    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
119        Ok(Self {
120            topics: decoder
121                .read_array(
122                    "offset fetch topic responses",
123                    OffsetFetchTopicResponse::decode,
124                )?
125                .unwrap_or_default(),
126            error_code: decoder.read_i16()?,
127        })
128    }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
132pub struct OffsetFetchResponseV9 {
133    pub throttle_time_ms: i32,
134    pub groups: Vec<OffsetFetchGroupResponse>,
135}
136
137impl OffsetFetchResponseV9 {
138    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
139        let throttle_time_ms = decoder.read_i32()?;
140        let groups = decoder
141            .read_compact_array("offset fetch group responses", |decoder| {
142                let group_id = decoder.read_compact_string()?;
143                let topics = decoder
144                    .read_compact_array("offset fetch topic responses", |decoder| {
145                        let name = decoder.read_compact_string()?;
146                        let partitions = decoder
147                            .read_compact_array("offset fetch partition responses", |decoder| {
148                                let partition_index = decoder.read_i32()?;
149                                let committed_offset = decoder.read_i64()?;
150                                let _committed_leader_epoch = decoder.read_i32()?;
151                                let metadata = decoder.read_compact_nullable_string()?;
152                                let error_code = decoder.read_i16()?;
153                                decoder.read_tagged_fields()?;
154                                Ok(OffsetFetchPartitionResponse {
155                                    partition_index,
156                                    committed_offset,
157                                    metadata,
158                                    error_code,
159                                })
160                            })?
161                            .unwrap_or_default();
162                        decoder.read_tagged_fields()?;
163                        Ok(OffsetFetchTopicResponse { name, partitions })
164                    })?
165                    .unwrap_or_default();
166                let error_code = decoder.read_i16()?;
167                decoder.read_tagged_fields()?;
168                Ok(OffsetFetchGroupResponse {
169                    group_id,
170                    topics,
171                    error_code,
172                })
173            })?
174            .unwrap_or_default();
175        decoder.read_tagged_fields()?;
176        Ok(Self {
177            throttle_time_ms,
178            groups,
179        })
180    }
181}
182
183#[derive(Debug, Clone, PartialEq, Eq)]
184pub struct OffsetFetchGroupResponse {
185    pub group_id: String,
186    pub topics: Vec<OffsetFetchTopicResponse>,
187    pub error_code: i16,
188}
189
190#[derive(Debug, Clone, PartialEq, Eq)]
191pub struct OffsetFetchTopicResponse {
192    pub name: String,
193    pub partitions: Vec<OffsetFetchPartitionResponse>,
194}
195
196impl OffsetFetchTopicResponse {
197    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
198        Ok(Self {
199            name: decoder.read_string()?,
200            partitions: decoder
201                .read_array(
202                    "offset fetch partition responses",
203                    OffsetFetchPartitionResponse::decode,
204                )?
205                .unwrap_or_default(),
206        })
207    }
208}
209
210#[derive(Debug, Clone, PartialEq, Eq)]
211pub struct OffsetFetchPartitionResponse {
212    pub partition_index: i32,
213    pub committed_offset: i64,
214    pub metadata: Option<String>,
215    pub error_code: i16,
216}
217
218impl OffsetFetchPartitionResponse {
219    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
220        Ok(Self {
221            partition_index: decoder.read_i32()?,
222            committed_offset: decoder.read_i64()?,
223            metadata: decoder.read_nullable_string()?,
224            error_code: decoder.read_i16()?,
225        })
226    }
227}
228
229#[cfg(test)]
230#[allow(clippy::unwrap_used)]
231mod tests {
232    use super::{
233        OffsetFetchGroupResponse, OffsetFetchPartitionResponse, OffsetFetchRequestV2,
234        OffsetFetchRequestV9, OffsetFetchResponseV2, OffsetFetchResponseV9, OffsetFetchTopic,
235        OffsetFetchTopicResponse, OffsetFetchTopicV9,
236    };
237    use crate::codec::{Decoder, Encoder};
238
239    #[test]
240    fn encodes_offset_fetch_v2_request_for_partitions() {
241        let request = OffsetFetchRequestV2 {
242            correlation_id: 29,
243            client_id: Some("kafrust".to_owned()),
244            group_id: "orders-group".to_owned(),
245            topics: Some(vec![OffsetFetchTopic {
246                name: "orders".to_owned(),
247                partition_indexes: vec![0, 1],
248            }]),
249        };
250
251        assert_eq!(
252            request.encode().unwrap(),
253            [
254                0, 9, // api key
255                0, 2, // api version
256                0, 0, 0, 29, // correlation id
257                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // client id
258                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
259                b'p', // group id
260                0, 0, 0, 1, // topic count
261                0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
262                0, 0, 0, 2, // partition count
263                0, 0, 0, 0, // partition 0
264                0, 0, 0, 1, // partition 1
265            ]
266        );
267    }
268
269    #[test]
270    fn encodes_offset_fetch_v2_request_for_all_topics() {
271        let request = OffsetFetchRequestV2 {
272            correlation_id: 31,
273            client_id: None,
274            group_id: "orders-group".to_owned(),
275            topics: None,
276        };
277
278        assert_eq!(
279            request.encode().unwrap(),
280            [
281                0, 9, // api key
282                0, 2, // api version
283                0, 0, 0, 31, // correlation id
284                0xff, 0xff, // null client id
285                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
286                b'p', // group id
287                0xff, 0xff, 0xff, 0xff, // null topics
288            ]
289        );
290    }
291
292    #[test]
293    fn decodes_offset_fetch_v2_response() {
294        let mut bytes = Encoder::new();
295        bytes.write_i32(1);
296        bytes.write_string("orders").unwrap();
297        bytes.write_i32(1);
298        bytes.write_i32(0);
299        bytes.write_i64(42);
300        bytes.write_nullable_string(Some("processed")).unwrap();
301        bytes.write_i16(0);
302        bytes.write_i16(0);
303        let bytes = bytes.into_bytes();
304
305        let mut decoder = Decoder::new(&bytes);
306        let response = OffsetFetchResponseV2::decode_body(&mut decoder).unwrap();
307
308        assert_eq!(
309            response,
310            OffsetFetchResponseV2 {
311                topics: vec![OffsetFetchTopicResponse {
312                    name: "orders".to_owned(),
313                    partitions: vec![OffsetFetchPartitionResponse {
314                        partition_index: 0,
315                        committed_offset: 42,
316                        metadata: Some("processed".to_owned()),
317                        error_code: 0,
318                    }],
319                }],
320                error_code: 0,
321            }
322        );
323        assert!(decoder.is_empty());
324    }
325
326    #[test]
327    fn encodes_offset_fetch_v9_request_for_consumer_protocol() {
328        let request = OffsetFetchRequestV9 {
329            correlation_id: 29,
330            client_id: Some("kafrust".to_owned()),
331            group_id: "orders-group".to_owned(),
332            member_id: Some("member-a".to_owned()),
333            member_epoch: 7,
334            topics: Some(vec![OffsetFetchTopicV9 {
335                name: "orders".to_owned(),
336                partition_indexes: vec![0, 1],
337            }]),
338            require_stable: false,
339        };
340
341        let encoded = request.encode().unwrap();
342        assert_eq!(&encoded[0..4], &[0, 9, 0, 9]);
343        let mut decoder = Decoder::new(&encoded[18..]);
344        let groups = decoder
345            .read_compact_array("offset fetch groups", |decoder| {
346                let group_id = decoder.read_compact_string()?;
347                let member_id = decoder.read_compact_nullable_string()?;
348                let member_epoch = decoder.read_i32()?;
349                let topics = decoder
350                    .read_compact_array("offset fetch topics", |decoder| {
351                        let name = decoder.read_compact_string()?;
352                        let partitions = decoder
353                            .read_compact_array("offset fetch partitions", |decoder| {
354                                decoder.read_i32()
355                            })?
356                            .unwrap_or_default();
357                        decoder.read_tagged_fields()?;
358                        Ok((name, partitions))
359                    })?
360                    .unwrap_or_default();
361                decoder.read_tagged_fields()?;
362                Ok((group_id, member_id, member_epoch, topics))
363            })
364            .unwrap()
365            .unwrap();
366        assert_eq!(groups[0].0, "orders-group");
367        assert_eq!(groups[0].1, Some("member-a".to_owned()));
368        assert_eq!(groups[0].2, 7);
369        assert_eq!(groups[0].3[0].0, "orders");
370        assert_eq!(groups[0].3[0].1, vec![0, 1]);
371        assert!(!decoder.read_bool().unwrap());
372        decoder.read_tagged_fields().unwrap();
373        assert!(decoder.is_empty());
374    }
375
376    #[test]
377    fn decodes_offset_fetch_v9_response() {
378        let mut bytes = Encoder::new();
379        bytes.write_i32(12);
380        bytes
381            .write_compact_array(Some(&[()]), |encoder, ()| {
382                encoder.write_compact_string("orders-group")?;
383                encoder.write_compact_array(Some(&[()]), |encoder, ()| {
384                    encoder.write_compact_string("orders")?;
385                    encoder.write_compact_array(Some(&[()]), |encoder, ()| {
386                        encoder.write_i32(0);
387                        encoder.write_i64(42);
388                        encoder.write_i32(-1);
389                        encoder.write_compact_nullable_string(Some("processed"))?;
390                        encoder.write_i16(0);
391                        encoder.write_empty_tagged_fields();
392                        Ok(())
393                    })?;
394                    encoder.write_empty_tagged_fields();
395                    Ok(())
396                })?;
397                encoder.write_i16(0);
398                encoder.write_empty_tagged_fields();
399                Ok(())
400            })
401            .unwrap();
402        bytes.write_empty_tagged_fields();
403
404        let bytes = bytes.into_bytes();
405        let mut decoder = Decoder::new(&bytes);
406        let response = OffsetFetchResponseV9::decode_body(&mut decoder).unwrap();
407
408        assert_eq!(response.throttle_time_ms, 12);
409        assert_eq!(
410            response.groups,
411            vec![OffsetFetchGroupResponse {
412                group_id: "orders-group".to_owned(),
413                topics: vec![OffsetFetchTopicResponse {
414                    name: "orders".to_owned(),
415                    partitions: vec![OffsetFetchPartitionResponse {
416                        partition_index: 0,
417                        committed_offset: 42,
418                        metadata: Some("processed".to_owned()),
419                        error_code: 0,
420                    }],
421                }],
422                error_code: 0,
423            }]
424        );
425        assert!(decoder.is_empty());
426    }
427}