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)]
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, 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',
259 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, ]
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, 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',
286 b'p', 0xff, 0xff, 0xff, 0xff, ]
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}