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#[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)]
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#[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, 0, 2, 0, 0, 0, 29, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', 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', 0, 0, 0, 1, 0, 6, b'o', b'r', b'd', b'e', b'r', b's', 0, 0, 0, 2, 0, 0, 0, 0, 0, 0, 0, 1, ]
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, 0, 2, 0, 0, 0, 31, 0xff, 0xff, 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', 0xff, 0xff, 0xff, 0xff, ]
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}