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/// OffsetFetch v10 request using topic UUIDs.
72#[derive(Debug, Clone, PartialEq, Eq)]
73pub struct OffsetFetchRequestV10 {
74    pub correlation_id: i32,
75    pub client_id: Option<String>,
76    pub group_id: String,
77    pub member_id: Option<String>,
78    pub member_epoch: i32,
79    pub topics: Option<Vec<OffsetFetchTopicV10>>,
80    pub require_stable: bool,
81}
82
83impl OffsetFetchRequestV10 {
84    pub fn encode(&self) -> Result<Vec<u8>> {
85        let mut encoder = Encoder::new();
86        RequestHeader {
87            api_key: API_KEY,
88            api_version: 10,
89            correlation_id: self.correlation_id,
90            client_id: self.client_id.clone(),
91        }
92        .encode_v2(&mut encoder)?;
93        encoder.write_compact_array(Some(&[()]), |encoder, ()| {
94            encoder.write_compact_string(&self.group_id)?;
95            encoder.write_compact_nullable_string(self.member_id.as_deref())?;
96            encoder.write_i32(self.member_epoch);
97            encoder.write_compact_array(self.topics.as_deref(), |encoder, topic| {
98                topic.encode(encoder)
99            })?;
100            encoder.write_empty_tagged_fields();
101            Ok(())
102        })?;
103        encoder.write_bool(self.require_stable);
104        encoder.write_empty_tagged_fields();
105        Ok(encoder.into_bytes())
106    }
107}
108
109#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct OffsetFetchTopic {
111    pub name: String,
112    pub partition_indexes: Vec<i32>,
113}
114
115impl OffsetFetchTopic {
116    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
117        encoder.write_string(&self.name)?;
118        encoder.write_array(
119            Some(self.partition_indexes.as_slice()),
120            |encoder, partition| {
121                encoder.write_i32(*partition);
122                Ok(())
123            },
124        )
125    }
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub struct OffsetFetchTopicV9 {
130    pub name: String,
131    pub partition_indexes: Vec<i32>,
132}
133
134impl OffsetFetchTopicV9 {
135    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
136        encoder.write_compact_string(&self.name)?;
137        encoder.write_compact_array(
138            Some(self.partition_indexes.as_slice()),
139            |encoder, partition| {
140                encoder.write_i32(*partition);
141                Ok(())
142            },
143        )?;
144        encoder.write_empty_tagged_fields();
145        Ok(())
146    }
147}
148
149#[derive(Debug, Clone, PartialEq, Eq)]
150pub struct OffsetFetchTopicV10 {
151    pub topic_id: [u8; 16],
152    pub partition_indexes: Vec<i32>,
153}
154
155impl OffsetFetchTopicV10 {
156    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
157        encoder.write_uuid(&self.topic_id);
158        encoder.write_compact_array(
159            Some(self.partition_indexes.as_slice()),
160            |encoder, partition| {
161                encoder.write_i32(*partition);
162                Ok(())
163            },
164        )?;
165        encoder.write_empty_tagged_fields();
166        Ok(())
167    }
168}
169
170#[derive(Debug, Clone, PartialEq, Eq)]
171pub struct OffsetFetchResponseV2 {
172    pub topics: Vec<OffsetFetchTopicResponse>,
173    pub error_code: i16,
174}
175
176impl OffsetFetchResponseV2 {
177    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
178        Ok(Self {
179            topics: decoder
180                .read_array(
181                    "offset fetch topic responses",
182                    OffsetFetchTopicResponse::decode,
183                )?
184                .unwrap_or_default(),
185            error_code: decoder.read_i16()?,
186        })
187    }
188}
189
190#[derive(Debug, Clone, PartialEq, Eq)]
191pub struct OffsetFetchResponseV9 {
192    pub throttle_time_ms: i32,
193    pub groups: Vec<OffsetFetchGroupResponse>,
194}
195
196impl OffsetFetchResponseV9 {
197    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
198        let throttle_time_ms = decoder.read_i32()?;
199        let groups = decoder
200            .read_compact_array("offset fetch group responses", |decoder| {
201                let group_id = decoder.read_compact_string()?;
202                let topics = decoder
203                    .read_compact_array("offset fetch topic responses", |decoder| {
204                        let name = decoder.read_compact_string()?;
205                        let partitions = decoder
206                            .read_compact_array("offset fetch partition responses", |decoder| {
207                                let partition_index = decoder.read_i32()?;
208                                let committed_offset = decoder.read_i64()?;
209                                let _committed_leader_epoch = decoder.read_i32()?;
210                                let metadata = decoder.read_compact_nullable_string()?;
211                                let error_code = decoder.read_i16()?;
212                                decoder.read_tagged_fields()?;
213                                Ok(OffsetFetchPartitionResponse {
214                                    partition_index,
215                                    committed_offset,
216                                    metadata,
217                                    error_code,
218                                })
219                            })?
220                            .unwrap_or_default();
221                        decoder.read_tagged_fields()?;
222                        Ok(OffsetFetchTopicResponse { name, partitions })
223                    })?
224                    .unwrap_or_default();
225                let error_code = decoder.read_i16()?;
226                decoder.read_tagged_fields()?;
227                Ok(OffsetFetchGroupResponse {
228                    group_id,
229                    topics,
230                    error_code,
231                })
232            })?
233            .unwrap_or_default();
234        decoder.read_tagged_fields()?;
235        Ok(Self {
236            throttle_time_ms,
237            groups,
238        })
239    }
240}
241
242/// OffsetFetch v10 response using topic UUIDs.
243#[derive(Debug, Clone, PartialEq, Eq)]
244pub struct OffsetFetchResponseV10 {
245    pub throttle_time_ms: i32,
246    pub groups: Vec<OffsetFetchGroupResponseV10>,
247}
248
249impl OffsetFetchResponseV10 {
250    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
251        let throttle_time_ms = decoder.read_i32()?;
252        let groups = decoder
253            .read_compact_array("offset fetch group UUID responses", |decoder| {
254                let group_id = decoder.read_compact_string()?;
255                let topics = decoder
256                    .read_compact_array("offset fetch topic UUID responses", |decoder| {
257                        let topic_id = decoder.read_uuid()?;
258                        let partitions = decoder
259                            .read_compact_array("offset fetch partition responses", |decoder| {
260                                let partition_index = decoder.read_i32()?;
261                                let committed_offset = decoder.read_i64()?;
262                                let _committed_leader_epoch = decoder.read_i32()?;
263                                let metadata = decoder.read_compact_nullable_string()?;
264                                let error_code = decoder.read_i16()?;
265                                decoder.read_tagged_fields()?;
266                                Ok(OffsetFetchPartitionResponse {
267                                    partition_index,
268                                    committed_offset,
269                                    metadata,
270                                    error_code,
271                                })
272                            })?
273                            .unwrap_or_default();
274                        decoder.read_tagged_fields()?;
275                        Ok(OffsetFetchTopicResponseV10 {
276                            topic_id,
277                            partitions,
278                        })
279                    })?
280                    .unwrap_or_default();
281                let error_code = decoder.read_i16()?;
282                decoder.read_tagged_fields()?;
283                Ok(OffsetFetchGroupResponseV10 {
284                    group_id,
285                    topics,
286                    error_code,
287                })
288            })?
289            .unwrap_or_default();
290        decoder.read_tagged_fields()?;
291        Ok(Self {
292            throttle_time_ms,
293            groups,
294        })
295    }
296}
297
298#[derive(Debug, Clone, PartialEq, Eq)]
299pub struct OffsetFetchGroupResponse {
300    pub group_id: String,
301    pub topics: Vec<OffsetFetchTopicResponse>,
302    pub error_code: i16,
303}
304
305#[derive(Debug, Clone, PartialEq, Eq)]
306pub struct OffsetFetchGroupResponseV10 {
307    pub group_id: String,
308    pub topics: Vec<OffsetFetchTopicResponseV10>,
309    pub error_code: i16,
310}
311
312#[derive(Debug, Clone, PartialEq, Eq)]
313pub struct OffsetFetchTopicResponse {
314    pub name: String,
315    pub partitions: Vec<OffsetFetchPartitionResponse>,
316}
317
318#[derive(Debug, Clone, PartialEq, Eq)]
319pub struct OffsetFetchTopicResponseV10 {
320    pub topic_id: [u8; 16],
321    pub partitions: Vec<OffsetFetchPartitionResponse>,
322}
323
324impl OffsetFetchTopicResponse {
325    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
326        Ok(Self {
327            name: decoder.read_string()?,
328            partitions: decoder
329                .read_array(
330                    "offset fetch partition responses",
331                    OffsetFetchPartitionResponse::decode,
332                )?
333                .unwrap_or_default(),
334        })
335    }
336}
337
338#[derive(Debug, Clone, PartialEq, Eq)]
339pub struct OffsetFetchPartitionResponse {
340    pub partition_index: i32,
341    pub committed_offset: i64,
342    pub metadata: Option<String>,
343    pub error_code: i16,
344}
345
346impl OffsetFetchPartitionResponse {
347    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
348        Ok(Self {
349            partition_index: decoder.read_i32()?,
350            committed_offset: decoder.read_i64()?,
351            metadata: decoder.read_nullable_string()?,
352            error_code: decoder.read_i16()?,
353        })
354    }
355}
356
357#[cfg(test)]
358#[allow(clippy::unwrap_used)]
359mod tests {
360    use super::{
361        OffsetFetchGroupResponse, OffsetFetchPartitionResponse, OffsetFetchRequestV10,
362        OffsetFetchRequestV2, OffsetFetchRequestV9, OffsetFetchResponseV10, OffsetFetchResponseV2,
363        OffsetFetchResponseV9, OffsetFetchTopic, OffsetFetchTopicResponse, OffsetFetchTopicV10,
364        OffsetFetchTopicV9,
365    };
366    use crate::codec::{Decoder, Encoder};
367
368    #[test]
369    fn encodes_offset_fetch_v2_request_for_partitions() {
370        let request = OffsetFetchRequestV2 {
371            correlation_id: 29,
372            client_id: Some("kafrust".to_owned()),
373            group_id: "orders-group".to_owned(),
374            topics: Some(vec![OffsetFetchTopic {
375                name: "orders".to_owned(),
376                partition_indexes: vec![0, 1],
377            }]),
378        };
379
380        assert_eq!(
381            request.encode().unwrap(),
382            [
383                0, 9, // api key
384                0, 2, // api version
385                0, 0, 0, 29, // correlation id
386                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // client id
387                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
388                b'p', // group id
389                0, 0, 0, 1, // topic count
390                0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
391                0, 0, 0, 2, // partition count
392                0, 0, 0, 0, // partition 0
393                0, 0, 0, 1, // partition 1
394            ]
395        );
396    }
397
398    #[test]
399    fn encodes_offset_fetch_v2_request_for_all_topics() {
400        let request = OffsetFetchRequestV2 {
401            correlation_id: 31,
402            client_id: None,
403            group_id: "orders-group".to_owned(),
404            topics: None,
405        };
406
407        assert_eq!(
408            request.encode().unwrap(),
409            [
410                0, 9, // api key
411                0, 2, // api version
412                0, 0, 0, 31, // correlation id
413                0xff, 0xff, // null client id
414                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
415                b'p', // group id
416                0xff, 0xff, 0xff, 0xff, // null topics
417            ]
418        );
419    }
420
421    #[test]
422    fn decodes_offset_fetch_v2_response() {
423        let mut bytes = Encoder::new();
424        bytes.write_i32(1);
425        bytes.write_string("orders").unwrap();
426        bytes.write_i32(1);
427        bytes.write_i32(0);
428        bytes.write_i64(42);
429        bytes.write_nullable_string(Some("processed")).unwrap();
430        bytes.write_i16(0);
431        bytes.write_i16(0);
432        let bytes = bytes.into_bytes();
433
434        let mut decoder = Decoder::new(&bytes);
435        let response = OffsetFetchResponseV2::decode_body(&mut decoder).unwrap();
436
437        assert_eq!(
438            response,
439            OffsetFetchResponseV2 {
440                topics: vec![OffsetFetchTopicResponse {
441                    name: "orders".to_owned(),
442                    partitions: vec![OffsetFetchPartitionResponse {
443                        partition_index: 0,
444                        committed_offset: 42,
445                        metadata: Some("processed".to_owned()),
446                        error_code: 0,
447                    }],
448                }],
449                error_code: 0,
450            }
451        );
452        assert!(decoder.is_empty());
453    }
454
455    #[test]
456    fn encodes_offset_fetch_v9_request_for_consumer_protocol() {
457        let request = OffsetFetchRequestV9 {
458            correlation_id: 29,
459            client_id: Some("kafrust".to_owned()),
460            group_id: "orders-group".to_owned(),
461            member_id: Some("member-a".to_owned()),
462            member_epoch: 7,
463            topics: Some(vec![OffsetFetchTopicV9 {
464                name: "orders".to_owned(),
465                partition_indexes: vec![0, 1],
466            }]),
467            require_stable: false,
468        };
469
470        let encoded = request.encode().unwrap();
471        assert_eq!(&encoded[0..4], &[0, 9, 0, 9]);
472        let mut decoder = Decoder::new(&encoded[18..]);
473        let groups = decoder
474            .read_compact_array("offset fetch groups", |decoder| {
475                let group_id = decoder.read_compact_string()?;
476                let member_id = decoder.read_compact_nullable_string()?;
477                let member_epoch = decoder.read_i32()?;
478                let topics = decoder
479                    .read_compact_array("offset fetch topics", |decoder| {
480                        let name = decoder.read_compact_string()?;
481                        let partitions = decoder
482                            .read_compact_array("offset fetch partitions", |decoder| {
483                                decoder.read_i32()
484                            })?
485                            .unwrap_or_default();
486                        decoder.read_tagged_fields()?;
487                        Ok((name, partitions))
488                    })?
489                    .unwrap_or_default();
490                decoder.read_tagged_fields()?;
491                Ok((group_id, member_id, member_epoch, topics))
492            })
493            .unwrap()
494            .unwrap();
495        assert_eq!(groups[0].0, "orders-group");
496        assert_eq!(groups[0].1, Some("member-a".to_owned()));
497        assert_eq!(groups[0].2, 7);
498        assert_eq!(groups[0].3[0].0, "orders");
499        assert_eq!(groups[0].3[0].1, vec![0, 1]);
500        assert!(!decoder.read_bool().unwrap());
501        decoder.read_tagged_fields().unwrap();
502        assert!(decoder.is_empty());
503    }
504
505    #[test]
506    fn decodes_offset_fetch_v9_response() {
507        let mut bytes = Encoder::new();
508        bytes.write_i32(12);
509        bytes
510            .write_compact_array(Some(&[()]), |encoder, ()| {
511                encoder.write_compact_string("orders-group")?;
512                encoder.write_compact_array(Some(&[()]), |encoder, ()| {
513                    encoder.write_compact_string("orders")?;
514                    encoder.write_compact_array(Some(&[()]), |encoder, ()| {
515                        encoder.write_i32(0);
516                        encoder.write_i64(42);
517                        encoder.write_i32(-1);
518                        encoder.write_compact_nullable_string(Some("processed"))?;
519                        encoder.write_i16(0);
520                        encoder.write_empty_tagged_fields();
521                        Ok(())
522                    })?;
523                    encoder.write_empty_tagged_fields();
524                    Ok(())
525                })?;
526                encoder.write_i16(0);
527                encoder.write_empty_tagged_fields();
528                Ok(())
529            })
530            .unwrap();
531        bytes.write_empty_tagged_fields();
532
533        let bytes = bytes.into_bytes();
534        let mut decoder = Decoder::new(&bytes);
535        let response = OffsetFetchResponseV9::decode_body(&mut decoder).unwrap();
536
537        assert_eq!(response.throttle_time_ms, 12);
538        assert_eq!(
539            response.groups,
540            vec![OffsetFetchGroupResponse {
541                group_id: "orders-group".to_owned(),
542                topics: vec![OffsetFetchTopicResponse {
543                    name: "orders".to_owned(),
544                    partitions: vec![OffsetFetchPartitionResponse {
545                        partition_index: 0,
546                        committed_offset: 42,
547                        metadata: Some("processed".to_owned()),
548                        error_code: 0,
549                    }],
550                }],
551                error_code: 0,
552            }]
553        );
554        assert!(decoder.is_empty());
555    }
556
557    #[test]
558    fn encodes_offset_fetch_v10_request_with_topic_uuid() {
559        let request = OffsetFetchRequestV10 {
560            correlation_id: 41,
561            client_id: Some("kafrust".to_owned()),
562            group_id: "orders-group".to_owned(),
563            member_id: Some("member-a".to_owned()),
564            member_epoch: 7,
565            topics: Some(vec![OffsetFetchTopicV10 {
566                topic_id: [9; 16],
567                partition_indexes: vec![0, 1],
568            }]),
569            require_stable: true,
570        };
571
572        let encoded = request.encode().unwrap();
573        assert_eq!(&encoded[0..4], &[0, 9, 0, 10]);
574        let mut decoder = Decoder::new(&encoded[18..]);
575        let groups = decoder
576            .read_compact_array("offset fetch groups", |decoder| {
577                let group_id = decoder.read_compact_string()?;
578                let member_id = decoder.read_compact_nullable_string()?;
579                let member_epoch = decoder.read_i32()?;
580                let topics = decoder
581                    .read_compact_array("offset fetch topics", |decoder| {
582                        let topic_id = decoder.read_uuid()?;
583                        let partitions = decoder
584                            .read_compact_array("offset fetch partitions", |decoder| {
585                                decoder.read_i32()
586                            })?
587                            .unwrap_or_default();
588                        decoder.read_tagged_fields()?;
589                        Ok((topic_id, partitions))
590                    })?
591                    .unwrap_or_default();
592                decoder.read_tagged_fields()?;
593                Ok((group_id, member_id, member_epoch, topics))
594            })
595            .unwrap()
596            .unwrap();
597        assert_eq!(groups[0].0, "orders-group");
598        assert_eq!(groups[0].1, Some("member-a".to_owned()));
599        assert_eq!(groups[0].2, 7);
600        assert_eq!(groups[0].3[0].0, [9; 16]);
601        assert_eq!(groups[0].3[0].1, vec![0, 1]);
602        assert!(decoder.read_bool().unwrap());
603        decoder.read_tagged_fields().unwrap();
604        assert!(decoder.is_empty());
605    }
606
607    #[test]
608    fn decodes_offset_fetch_v10_response_with_topic_uuid() {
609        let mut bytes = Encoder::new();
610        bytes.write_i32(12);
611        bytes
612            .write_compact_array(Some(&[()]), |encoder, ()| {
613                encoder.write_compact_string("orders-group")?;
614                encoder.write_compact_array(Some(&[()]), |encoder, ()| {
615                    encoder.write_uuid(&[10; 16]);
616                    encoder.write_compact_array(Some(&[()]), |encoder, ()| {
617                        encoder.write_i32(0);
618                        encoder.write_i64(42);
619                        encoder.write_i32(9);
620                        encoder.write_compact_nullable_string(Some("processed"))?;
621                        encoder.write_i16(0);
622                        encoder.write_empty_tagged_fields();
623                        Ok(())
624                    })?;
625                    encoder.write_empty_tagged_fields();
626                    Ok(())
627                })?;
628                encoder.write_i16(0);
629                encoder.write_empty_tagged_fields();
630                Ok(())
631            })
632            .unwrap();
633        bytes.write_empty_tagged_fields();
634
635        let bytes = bytes.into_bytes();
636        let mut decoder = Decoder::new(&bytes);
637        let response = OffsetFetchResponseV10::decode_body(&mut decoder).unwrap();
638        assert_eq!(response.throttle_time_ms, 12);
639        assert_eq!(response.groups[0].group_id, "orders-group");
640        assert_eq!(response.groups[0].topics[0].topic_id, [10; 16]);
641        assert_eq!(
642            response.groups[0].topics[0].partitions[0],
643            OffsetFetchPartitionResponse {
644                partition_index: 0,
645                committed_offset: 42,
646                metadata: Some("processed".to_owned()),
647                error_code: 0,
648            }
649        );
650        assert!(decoder.is_empty());
651    }
652}