Skip to main content

kafrust_protocol/api/
produce.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::{Error, Result};
3use crate::header::RequestHeader;
4use crate::record_batch::{compress_record_batch_records, RecordBatchCompression};
5
6pub const API_KEY: i16 = 0;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct ProduceRequestV2 {
10    pub correlation_id: i32,
11    pub client_id: Option<String>,
12    pub acks: i16,
13    pub timeout_ms: i32,
14    pub topics: Vec<ProduceTopicV2>,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct ProduceRequestV3 {
19    pub correlation_id: i32,
20    pub client_id: Option<String>,
21    pub transactional_id: Option<String>,
22    pub acks: i16,
23    pub timeout_ms: i32,
24    pub topics: Vec<ProduceTopicV3>,
25}
26
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct ProduceRequestV7 {
29    pub correlation_id: i32,
30    pub client_id: Option<String>,
31    pub transactional_id: Option<String>,
32    pub acks: i16,
33    pub timeout_ms: i32,
34    pub topics: Vec<ProduceTopicV3>,
35}
36
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct ProduceRequestV9 {
39    pub correlation_id: i32,
40    pub client_id: Option<String>,
41    pub transactional_id: Option<String>,
42    pub acks: i16,
43    pub timeout_ms: i32,
44    pub topics: Vec<ProduceTopicV3>,
45}
46
47/// Flexible Produce request used by Kafka API versions 11 and newer.
48///
49/// Kafka keeps the v9 flexible RecordBatch schema for these versions; the
50/// request header carries the negotiated API version.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct ProduceRequestV11 {
53    pub correlation_id: i32,
54    pub client_id: Option<String>,
55    pub transactional_id: Option<String>,
56    pub acks: i16,
57    pub timeout_ms: i32,
58    pub topics: Vec<ProduceTopicV3>,
59}
60
61/// Flexible Produce request for Kafka API version 12.
62///
63/// Kafka keeps the v9 flexible RecordBatch schema for this version; the
64/// request header carries the negotiated API version.
65#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct ProduceRequestV12 {
67    pub correlation_id: i32,
68    pub client_id: Option<String>,
69    pub transactional_id: Option<String>,
70    pub acks: i16,
71    pub timeout_ms: i32,
72    pub topics: Vec<ProduceTopicV3>,
73}
74
75/// Flexible Produce request for Kafka API version 13.
76///
77/// Kafka v13 replaces the topic name with the stable topic UUID while keeping
78/// the flexible RecordBatch partition schema from v9-v12.
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub struct ProduceRequestV13 {
81    pub correlation_id: i32,
82    pub client_id: Option<String>,
83    pub transactional_id: Option<String>,
84    pub acks: i16,
85    pub timeout_ms: i32,
86    pub topics: Vec<ProduceTopicV13>,
87}
88
89impl ProduceRequestV3 {
90    pub fn encode(&self) -> Result<Vec<u8>> {
91        encode_record_batch_request(
92            3,
93            self.correlation_id,
94            self.client_id.clone(),
95            self.transactional_id.as_deref(),
96            self.acks,
97            self.timeout_ms,
98            &self.topics,
99        )
100    }
101}
102
103impl ProduceRequestV7 {
104    pub fn encode(&self) -> Result<Vec<u8>> {
105        encode_record_batch_request(
106            7,
107            self.correlation_id,
108            self.client_id.clone(),
109            self.transactional_id.as_deref(),
110            self.acks,
111            self.timeout_ms,
112            &self.topics,
113        )
114    }
115}
116
117impl ProduceRequestV9 {
118    pub fn encode(&self) -> Result<Vec<u8>> {
119        encode_flexible_produce_request(
120            9,
121            self.correlation_id,
122            self.client_id.clone(),
123            self.transactional_id.as_deref(),
124            self.acks,
125            self.timeout_ms,
126            &self.topics,
127        )
128    }
129}
130
131impl ProduceRequestV11 {
132    pub fn encode(&self) -> Result<Vec<u8>> {
133        encode_flexible_produce_request(
134            11,
135            self.correlation_id,
136            self.client_id.clone(),
137            self.transactional_id.as_deref(),
138            self.acks,
139            self.timeout_ms,
140            &self.topics,
141        )
142    }
143}
144
145impl ProduceRequestV12 {
146    pub fn encode(&self) -> Result<Vec<u8>> {
147        encode_flexible_produce_request(
148            12,
149            self.correlation_id,
150            self.client_id.clone(),
151            self.transactional_id.as_deref(),
152            self.acks,
153            self.timeout_ms,
154            &self.topics,
155        )
156    }
157}
158
159impl ProduceRequestV13 {
160    pub fn encode(&self) -> Result<Vec<u8>> {
161        let mut encoder = Encoder::new();
162        RequestHeader {
163            api_key: API_KEY,
164            api_version: 13,
165            correlation_id: self.correlation_id,
166            client_id: self.client_id.clone(),
167        }
168        .encode_v2(&mut encoder)?;
169        encoder.write_compact_nullable_string(self.transactional_id.as_deref())?;
170        encoder.write_i16(self.acks);
171        encoder.write_i32(self.timeout_ms);
172        encoder.write_compact_array(Some(self.topics.as_slice()), |encoder, topic| {
173            topic.encode(encoder, self.transactional_id.is_some())
174        })?;
175        encoder.write_empty_tagged_fields();
176        Ok(encoder.into_bytes())
177    }
178}
179
180fn encode_flexible_produce_request(
181    api_version: i16,
182    correlation_id: i32,
183    client_id: Option<String>,
184    transactional_id: Option<&str>,
185    acks: i16,
186    timeout_ms: i32,
187    topics: &[ProduceTopicV3],
188) -> Result<Vec<u8>> {
189    let mut encoder = Encoder::new();
190    RequestHeader {
191        api_key: API_KEY,
192        api_version,
193        correlation_id,
194        client_id,
195    }
196    .encode_v2(&mut encoder)?;
197    encoder.write_compact_nullable_string(transactional_id)?;
198    encoder.write_i16(acks);
199    encoder.write_i32(timeout_ms);
200    encoder.write_compact_array(Some(topics), |encoder, topic| {
201        topic.encode_v9(encoder, transactional_id.is_some())
202    })?;
203    encoder.write_empty_tagged_fields();
204    Ok(encoder.into_bytes())
205}
206
207fn encode_record_batch_request(
208    api_version: i16,
209    correlation_id: i32,
210    client_id: Option<String>,
211    transactional_id: Option<&str>,
212    acks: i16,
213    timeout_ms: i32,
214    topics: &[ProduceTopicV3],
215) -> Result<Vec<u8>> {
216    let mut encoder = Encoder::new();
217    RequestHeader {
218        api_key: API_KEY,
219        api_version,
220        correlation_id,
221        client_id,
222    }
223    .encode_v1(&mut encoder)?;
224    encoder.write_nullable_string(transactional_id)?;
225    encoder.write_i16(acks);
226    encoder.write_i32(timeout_ms);
227    encoder.write_array(Some(topics), |encoder, topic| {
228        topic.encode(encoder, transactional_id.is_some())
229    })?;
230    Ok(encoder.into_bytes())
231}
232
233impl ProduceRequestV2 {
234    pub fn encode(&self) -> Result<Vec<u8>> {
235        let mut encoder = Encoder::new();
236        RequestHeader {
237            api_key: API_KEY,
238            api_version: 2,
239            correlation_id: self.correlation_id,
240            client_id: self.client_id.clone(),
241        }
242        .encode_v1(&mut encoder)?;
243        encoder.write_i16(self.acks);
244        encoder.write_i32(self.timeout_ms);
245        encoder.write_array(Some(self.topics.as_slice()), |encoder, topic| {
246            topic.encode(encoder)
247        })?;
248        Ok(encoder.into_bytes())
249    }
250}
251
252#[derive(Debug, Clone, PartialEq, Eq)]
253pub struct ProduceTopicV2 {
254    pub name: String,
255    pub partitions: Vec<ProducePartitionV2>,
256}
257
258#[derive(Debug, Clone, PartialEq, Eq)]
259pub struct ProduceTopicV3 {
260    pub name: String,
261    pub partitions: Vec<ProducePartitionV3>,
262}
263
264/// Produce topic payload for Kafka API version 13.
265#[derive(Debug, Clone, PartialEq, Eq)]
266pub struct ProduceTopicV13 {
267    pub topic_id: [u8; 16],
268    pub partitions: Vec<ProducePartitionV3>,
269}
270
271impl ProduceTopicV13 {
272    fn encode(&self, encoder: &mut Encoder, transactional: bool) -> Result<()> {
273        encoder.write_uuid(&self.topic_id);
274        encoder.write_compact_array(Some(self.partitions.as_slice()), |encoder, partition| {
275            partition.encode_v9(encoder, transactional)
276        })?;
277        encoder.write_empty_tagged_fields();
278        Ok(())
279    }
280}
281
282impl ProduceTopicV3 {
283    fn encode(&self, encoder: &mut Encoder, transactional: bool) -> Result<()> {
284        encoder.write_string(&self.name)?;
285        encoder.write_array(Some(self.partitions.as_slice()), |encoder, partition| {
286            partition.encode(encoder, transactional)
287        })
288    }
289
290    fn encode_v9(&self, encoder: &mut Encoder, transactional: bool) -> Result<()> {
291        encoder.write_compact_string(&self.name)?;
292        encoder.write_compact_array(Some(self.partitions.as_slice()), |encoder, partition| {
293            partition.encode_v9(encoder, transactional)
294        })?;
295        encoder.write_empty_tagged_fields();
296        Ok(())
297    }
298}
299
300impl ProduceTopicV2 {
301    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
302        encoder.write_string(&self.name)?;
303        encoder.write_array(Some(self.partitions.as_slice()), |encoder, partition| {
304            partition.encode(encoder)
305        })
306    }
307}
308
309#[derive(Debug, Clone, PartialEq, Eq)]
310pub struct ProducePartitionV2 {
311    pub partition_index: i32,
312    pub records: Vec<MessageSetMessage>,
313}
314
315#[derive(Debug, Clone, PartialEq, Eq)]
316pub struct ProducePartitionV3 {
317    pub partition_index: i32,
318    pub records: Vec<RecordBatchMessage>,
319    pub compression: RecordBatchCompression,
320    pub identity: RecordBatchIdentity,
321}
322
323#[derive(Debug, Clone, Copy, PartialEq, Eq)]
324pub struct RecordBatchIdentity {
325    pub producer_id: i64,
326    pub producer_epoch: i16,
327    pub base_sequence: i32,
328}
329
330impl RecordBatchIdentity {
331    pub const NON_IDEMPOTENT: Self = Self {
332        producer_id: -1,
333        producer_epoch: -1,
334        base_sequence: -1,
335    };
336}
337
338impl ProducePartitionV3 {
339    fn encode(&self, encoder: &mut Encoder, transactional: bool) -> Result<()> {
340        encoder.write_i32(self.partition_index);
341        let record_set = encode_record_batch_set_with_compression_identity_and_transaction(
342            &self.records,
343            self.compression,
344            self.identity,
345            transactional,
346        )?;
347        encoder.write_bytes(&record_set)
348    }
349
350    fn encode_v9(&self, encoder: &mut Encoder, transactional: bool) -> Result<()> {
351        encoder.write_i32(self.partition_index);
352        let record_set = encode_record_batch_set_with_compression_identity_and_transaction(
353            &self.records,
354            self.compression,
355            self.identity,
356            transactional,
357        )?;
358        encoder.write_compact_nullable_bytes(Some(&record_set))?;
359        encoder.write_empty_tagged_fields();
360        Ok(())
361    }
362}
363
364impl ProducePartitionV2 {
365    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
366        encoder.write_i32(self.partition_index);
367        let record_set = encode_message_set(&self.records)?;
368        encoder.write_bytes(&record_set)
369    }
370}
371
372/// Returns the encoded byte length of a Produce v2 message set.
373pub fn encoded_message_set_len(records: &[MessageSetMessage]) -> Result<usize> {
374    Ok(encode_message_set(records)?.len())
375}
376
377/// Returns the encoded byte length of a Produce v3 record batch set.
378pub fn encoded_record_batch_set_len(records: &[RecordBatchMessage]) -> Result<usize> {
379    Ok(encode_record_batch_set(records)?.len())
380}
381
382/// Returns the encoded byte length of a Produce v3 record batch set.
383pub fn encoded_record_batch_set_len_with_compression(
384    records: &[RecordBatchMessage],
385    compression: RecordBatchCompression,
386) -> Result<usize> {
387    encoded_record_batch_set_len_with_compression_and_identity(
388        records,
389        compression,
390        RecordBatchIdentity::NON_IDEMPOTENT,
391    )
392}
393
394pub fn encoded_record_batch_set_len_with_compression_and_identity(
395    records: &[RecordBatchMessage],
396    compression: RecordBatchCompression,
397    identity: RecordBatchIdentity,
398) -> Result<usize> {
399    Ok(
400        encode_record_batch_set_with_compression_and_identity(records, compression, identity)?
401            .len(),
402    )
403}
404
405#[derive(Debug, Clone, PartialEq, Eq)]
406pub struct MessageSetMessage {
407    pub key: Option<Vec<u8>>,
408    pub value: Option<Vec<u8>>,
409    pub timestamp_ms: i64,
410}
411
412impl MessageSetMessage {
413    pub fn new(key: Option<Vec<u8>>, value: Option<Vec<u8>>, timestamp_ms: i64) -> Self {
414        Self {
415            key,
416            value,
417            timestamp_ms,
418        }
419    }
420}
421
422#[derive(Debug, Clone, PartialEq, Eq)]
423pub struct RecordBatchHeader {
424    pub key: String,
425    pub value: Option<Vec<u8>>,
426}
427
428impl RecordBatchHeader {
429    pub fn new(key: impl Into<String>, value: Option<Vec<u8>>) -> Self {
430        Self {
431            key: key.into(),
432            value,
433        }
434    }
435}
436
437#[derive(Debug, Clone, PartialEq, Eq)]
438pub struct RecordBatchMessage {
439    pub key: Option<Vec<u8>>,
440    pub value: Option<Vec<u8>>,
441    pub timestamp_ms: i64,
442    pub headers: Vec<RecordBatchHeader>,
443}
444
445impl RecordBatchMessage {
446    pub fn new(key: Option<Vec<u8>>, value: Option<Vec<u8>>, timestamp_ms: i64) -> Self {
447        Self {
448            key,
449            value,
450            timestamp_ms,
451            headers: Vec::new(),
452        }
453    }
454
455    pub fn header(mut self, key: impl Into<String>, value: Option<Vec<u8>>) -> Self {
456        self.headers.push(RecordBatchHeader::new(key, value));
457        self
458    }
459}
460
461#[derive(Debug, Clone, PartialEq, Eq)]
462pub struct ProduceResponseV2 {
463    pub responses: Vec<ProduceTopicResponseV2>,
464    pub throttle_time_ms: i32,
465}
466
467impl ProduceResponseV2 {
468    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
469        Ok(Self {
470            responses: decoder
471                .read_array("produce responses", ProduceTopicResponseV2::decode)?
472                .unwrap_or_default(),
473            throttle_time_ms: decoder.read_i32()?,
474        })
475    }
476}
477
478#[derive(Debug, Clone, PartialEq, Eq)]
479pub struct ProduceTopicResponseV2 {
480    pub name: String,
481    pub partitions: Vec<ProducePartitionResponseV2>,
482}
483
484impl ProduceTopicResponseV2 {
485    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
486        Ok(Self {
487            name: decoder.read_string()?,
488            partitions: decoder
489                .read_array(
490                    "produce partition responses",
491                    ProducePartitionResponseV2::decode,
492                )?
493                .unwrap_or_default(),
494        })
495    }
496}
497
498#[derive(Debug, Clone, PartialEq, Eq)]
499pub struct ProducePartitionResponseV2 {
500    pub partition_index: i32,
501    pub error_code: i16,
502    pub base_offset: i64,
503    pub log_append_time_ms: i64,
504}
505
506impl ProducePartitionResponseV2 {
507    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
508        Ok(Self {
509            partition_index: decoder.read_i32()?,
510            error_code: decoder.read_i16()?,
511            base_offset: decoder.read_i64()?,
512            log_append_time_ms: decoder.read_i64()?,
513        })
514    }
515}
516
517#[derive(Debug, Clone, PartialEq, Eq)]
518pub struct ProduceResponseV7 {
519    pub responses: Vec<ProduceTopicResponseV7>,
520    pub throttle_time_ms: i32,
521}
522
523#[derive(Debug, Clone, PartialEq, Eq)]
524pub struct ProduceResponseV9 {
525    pub responses: Vec<ProduceTopicResponseV9>,
526    pub throttle_time_ms: i32,
527}
528
529/// Flexible Produce response used by Kafka API versions 11 and newer.
530pub type ProduceResponseV11 = ProduceResponseV9;
531
532/// Flexible Produce response for Kafka API version 12.
533pub type ProduceResponseV12 = ProduceResponseV9;
534
535/// Flexible Produce response for Kafka API version 13.
536#[derive(Debug, Clone, PartialEq, Eq)]
537pub struct ProduceResponseV13 {
538    pub responses: Vec<ProduceTopicResponseV13>,
539    pub throttle_time_ms: i32,
540}
541
542impl ProduceResponseV13 {
543    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
544        let responses = decoder
545            .read_compact_array("produce responses", ProduceTopicResponseV13::decode)?
546            .unwrap_or_default();
547        let throttle_time_ms = decoder.read_i32()?;
548        decoder.read_tagged_fields()?;
549        Ok(Self {
550            responses,
551            throttle_time_ms,
552        })
553    }
554}
555
556#[derive(Debug, Clone, PartialEq, Eq)]
557pub struct ProduceTopicResponseV13 {
558    pub topic_id: [u8; 16],
559    pub partitions: Vec<ProducePartitionResponseV9>,
560}
561
562impl ProduceTopicResponseV13 {
563    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
564        let topic_id = decoder.read_uuid()?;
565        let partitions = decoder
566            .read_compact_array(
567                "produce partition responses",
568                ProducePartitionResponseV9::decode,
569            )?
570            .unwrap_or_default();
571        decoder.read_tagged_fields()?;
572        Ok(Self {
573            topic_id,
574            partitions,
575        })
576    }
577}
578
579impl ProduceResponseV9 {
580    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
581        let responses = decoder
582            .read_compact_array("produce responses", ProduceTopicResponseV9::decode)?
583            .unwrap_or_default();
584        let throttle_time_ms = decoder.read_i32()?;
585        decoder.read_tagged_fields()?;
586        Ok(Self {
587            responses,
588            throttle_time_ms,
589        })
590    }
591}
592
593#[derive(Debug, Clone, PartialEq, Eq)]
594pub struct ProduceTopicResponseV9 {
595    pub name: String,
596    pub partitions: Vec<ProducePartitionResponseV9>,
597}
598
599impl ProduceTopicResponseV9 {
600    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
601        let name = decoder.read_compact_string()?;
602        let partitions = decoder
603            .read_compact_array(
604                "produce partition responses",
605                ProducePartitionResponseV9::decode,
606            )?
607            .unwrap_or_default();
608        decoder.read_tagged_fields()?;
609        Ok(Self { name, partitions })
610    }
611}
612
613#[derive(Debug, Clone, PartialEq, Eq)]
614pub struct ProducePartitionResponseV9 {
615    pub partition_index: i32,
616    pub error_code: i16,
617    pub base_offset: i64,
618    pub log_append_time_ms: i64,
619    pub log_start_offset: i64,
620    pub record_errors: Vec<ProduceRecordErrorV9>,
621    pub error_message: Option<String>,
622}
623
624impl ProducePartitionResponseV9 {
625    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
626        let partition_index = decoder.read_i32()?;
627        let error_code = decoder.read_i16()?;
628        let base_offset = decoder.read_i64()?;
629        let log_append_time_ms = decoder.read_i64()?;
630        let log_start_offset = decoder.read_i64()?;
631        let record_errors = decoder
632            .read_compact_array("produce record errors", ProduceRecordErrorV9::decode)?
633            .unwrap_or_default();
634        let error_message = decoder.read_compact_nullable_string()?;
635        decoder.read_tagged_fields()?;
636        Ok(Self {
637            partition_index,
638            error_code,
639            base_offset,
640            log_append_time_ms,
641            log_start_offset,
642            record_errors,
643            error_message,
644        })
645    }
646}
647
648#[derive(Debug, Clone, PartialEq, Eq)]
649pub struct ProduceRecordErrorV9 {
650    pub batch_index: i32,
651    pub error_message: Option<String>,
652}
653
654impl ProduceRecordErrorV9 {
655    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
656        let batch_index = decoder.read_i32()?;
657        let error_message = decoder.read_compact_nullable_string()?;
658        decoder.read_tagged_fields()?;
659        Ok(Self {
660            batch_index,
661            error_message,
662        })
663    }
664}
665
666impl ProduceResponseV7 {
667    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
668        Ok(Self {
669            responses: decoder
670                .read_array("produce responses", ProduceTopicResponseV7::decode)?
671                .unwrap_or_default(),
672            throttle_time_ms: decoder.read_i32()?,
673        })
674    }
675}
676
677#[derive(Debug, Clone, PartialEq, Eq)]
678pub struct ProduceTopicResponseV7 {
679    pub name: String,
680    pub partitions: Vec<ProducePartitionResponseV7>,
681}
682
683impl ProduceTopicResponseV7 {
684    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
685        Ok(Self {
686            name: decoder.read_string()?,
687            partitions: decoder
688                .read_array(
689                    "produce partition responses",
690                    ProducePartitionResponseV7::decode,
691                )?
692                .unwrap_or_default(),
693        })
694    }
695}
696
697#[derive(Debug, Clone, PartialEq, Eq)]
698pub struct ProducePartitionResponseV7 {
699    pub partition_index: i32,
700    pub error_code: i16,
701    pub base_offset: i64,
702    pub log_append_time_ms: i64,
703    pub log_start_offset: i64,
704}
705
706impl ProducePartitionResponseV7 {
707    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
708        Ok(Self {
709            partition_index: decoder.read_i32()?,
710            error_code: decoder.read_i16()?,
711            base_offset: decoder.read_i64()?,
712            log_append_time_ms: decoder.read_i64()?,
713            log_start_offset: decoder.read_i64()?,
714        })
715    }
716}
717
718fn encode_message_set(records: &[MessageSetMessage]) -> Result<Vec<u8>> {
719    let mut set = Encoder::new();
720    for record in records {
721        let message = encode_message(record)?;
722        set.write_i64(0);
723        set.write_i32(i32::try_from(message.len()).map_err(|_| Error::LengthOverflow("message"))?);
724        set.write_raw(&message);
725    }
726    Ok(set.into_bytes())
727}
728
729fn encode_message(record: &MessageSetMessage) -> Result<Vec<u8>> {
730    let mut body = Encoder::new();
731    body.write_i8(1);
732    body.write_i8(0);
733    body.write_i64(record.timestamp_ms);
734    body.write_nullable_bytes(record.key.as_deref())?;
735    body.write_nullable_bytes(record.value.as_deref())?;
736    let body = body.into_bytes();
737
738    let mut message = Encoder::new();
739    message.write_i32(crc32_ieee(&body) as i32);
740    message.write_raw(&body);
741    Ok(message.into_bytes())
742}
743
744fn encode_record_batch_set(records: &[RecordBatchMessage]) -> Result<Vec<u8>> {
745    encode_record_batch_set_with_compression(records, RecordBatchCompression::None)
746}
747
748fn encode_record_batch_set_with_compression(
749    records: &[RecordBatchMessage],
750    compression: RecordBatchCompression,
751) -> Result<Vec<u8>> {
752    encode_record_batch_set_with_compression_and_identity(
753        records,
754        compression,
755        RecordBatchIdentity::NON_IDEMPOTENT,
756    )
757}
758
759fn encode_record_batch_set_with_compression_and_identity(
760    records: &[RecordBatchMessage],
761    compression: RecordBatchCompression,
762    identity: RecordBatchIdentity,
763) -> Result<Vec<u8>> {
764    encode_record_batch_set_with_compression_identity_and_transaction(
765        records,
766        compression,
767        identity,
768        false,
769    )
770}
771
772fn encode_record_batch_set_with_compression_identity_and_transaction(
773    records: &[RecordBatchMessage],
774    compression: RecordBatchCompression,
775    identity: RecordBatchIdentity,
776    transactional: bool,
777) -> Result<Vec<u8>> {
778    let base_timestamp = records
779        .first()
780        .map(|record| record.timestamp_ms)
781        .unwrap_or_default();
782    let max_timestamp = records
783        .iter()
784        .map(|record| record.timestamp_ms)
785        .max()
786        .unwrap_or(base_timestamp);
787    let last_offset_delta = records
788        .len()
789        .checked_sub(1)
790        .map(|delta| i32::try_from(delta).map_err(|_| Error::LengthOverflow("record batch")))
791        .transpose()?
792        .unwrap_or_default();
793
794    let record_count =
795        i32::try_from(records.len()).map_err(|_| Error::LengthOverflow("record batch records"))?;
796    let mut record_bytes = Encoder::new();
797    for (offset_delta, record) in records.iter().enumerate() {
798        let encoded = encode_record(record, base_timestamp, offset_delta)?;
799        record_bytes.write_varint(
800            i32::try_from(encoded.len()).map_err(|_| Error::LengthOverflow("record"))?,
801        );
802        record_bytes.write_raw(&encoded);
803    }
804    let record_bytes = compress_record_batch_records(compression, &record_bytes.into_bytes())?;
805
806    let mut crc_payload = Encoder::new();
807    let attributes = compression.attributes() | if transactional { 0x10 } else { 0 };
808    crc_payload.write_i16(attributes);
809    crc_payload.write_i32(last_offset_delta);
810    crc_payload.write_i64(base_timestamp);
811    crc_payload.write_i64(max_timestamp);
812    crc_payload.write_i64(identity.producer_id);
813    crc_payload.write_i16(identity.producer_epoch);
814    crc_payload.write_i32(identity.base_sequence);
815    crc_payload.write_i32(record_count);
816    crc_payload.write_raw(&record_bytes);
817    let crc_payload = crc_payload.into_bytes();
818
819    let mut batch = Encoder::new();
820    batch.write_i32(0);
821    batch.write_i8(2);
822    batch.write_i32(crc32c(&crc_payload) as i32);
823    batch.write_raw(&crc_payload);
824    let batch = batch.into_bytes();
825
826    let mut set = Encoder::new();
827    set.write_i64(0);
828    set.write_i32(i32::try_from(batch.len()).map_err(|_| Error::LengthOverflow("record batch"))?);
829    set.write_raw(&batch);
830    Ok(set.into_bytes())
831}
832
833fn encode_record(
834    record: &RecordBatchMessage,
835    base_timestamp: i64,
836    offset_delta: usize,
837) -> Result<Vec<u8>> {
838    let mut encoder = Encoder::new();
839    encoder.write_i8(0);
840    encoder.write_varlong(record.timestamp_ms.saturating_sub(base_timestamp));
841    encoder.write_varint(
842        i32::try_from(offset_delta).map_err(|_| Error::LengthOverflow("record offset delta"))?,
843    );
844    encoder.write_varint_nullable_bytes(record.key.as_deref())?;
845    encoder.write_varint_nullable_bytes(record.value.as_deref())?;
846    encoder.write_varint(
847        i32::try_from(record.headers.len()).map_err(|_| Error::LengthOverflow("record headers"))?,
848    );
849    for header in &record.headers {
850        encoder.write_varint_bytes(header.key.as_bytes())?;
851        encoder.write_varint_nullable_bytes(header.value.as_deref())?;
852    }
853    Ok(encoder.into_bytes())
854}
855
856fn crc32_ieee(bytes: &[u8]) -> u32 {
857    crc32_with_table(bytes, &CRC32_IEEE_TABLE)
858}
859
860fn crc32c(bytes: &[u8]) -> u32 {
861    crc32_with_table(bytes, &CRC32C_TABLE)
862}
863
864const CRC32_IEEE_TABLE: [u32; 256] = crc32_table(0xedb8_8320);
865const CRC32C_TABLE: [u32; 256] = crc32_table(0x82f6_3b78);
866
867const fn crc32_table(polynomial: u32) -> [u32; 256] {
868    let mut table = [0; 256];
869    let mut index = 0;
870    while index < table.len() {
871        let mut value = index as u32;
872        let mut bit = 0;
873        while bit < 8 {
874            let mask = 0u32.wrapping_sub(value & 1);
875            value = (value >> 1) ^ (polynomial & mask);
876            bit += 1;
877        }
878        table[index] = value;
879        index += 1;
880    }
881    table
882}
883
884fn crc32_with_table(bytes: &[u8], table: &[u32; 256]) -> u32 {
885    let mut crc = 0xffff_ffffu32;
886    for byte in bytes {
887        let index = usize::from((crc as u8) ^ byte);
888        crc = (crc >> 8) ^ table[index];
889    }
890    !crc
891}
892
893#[cfg(test)]
894#[allow(clippy::unwrap_used)]
895mod tests {
896    use super::{
897        crc32_ieee, crc32c, encode_message_set, encode_record_batch_set,
898        encode_record_batch_set_with_compression,
899        encode_record_batch_set_with_compression_and_identity,
900        encode_record_batch_set_with_compression_identity_and_transaction, encoded_message_set_len,
901        encoded_record_batch_set_len, MessageSetMessage, ProducePartitionV2, ProducePartitionV3,
902        ProduceRequestV11, ProduceRequestV12, ProduceRequestV13, ProduceRequestV2,
903        ProduceRequestV3, ProduceRequestV7, ProduceRequestV9, ProduceResponseV13,
904        ProduceResponseV2, ProduceResponseV7, ProduceResponseV9, ProduceTopicV13, ProduceTopicV2,
905        ProduceTopicV3, RecordBatchIdentity, RecordBatchMessage,
906    };
907    use crate::codec::{DecodeLimits, Decoder};
908    use crate::record_batch::RecordBatchCompression;
909    use crate::{api::fetch::FetchResponseV2, codec::Encoder};
910
911    #[test]
912    fn crc_implementations_match_standard_check_vectors() {
913        assert_eq!(crc32_ieee(b"123456789"), 0xcbf4_3926);
914        assert_eq!(crc32c(b"123456789"), 0xe306_9283);
915    }
916
917    #[test]
918    fn encodes_produce_request_v2() {
919        let request = ProduceRequestV2 {
920            correlation_id: 5,
921            client_id: Some("kafrust".to_owned()),
922            acks: 1,
923            timeout_ms: 30_000,
924            topics: vec![ProduceTopicV2 {
925                name: "orders".to_owned(),
926                partitions: vec![ProducePartitionV2 {
927                    partition_index: 0,
928                    records: vec![MessageSetMessage::new(
929                        Some(b"order-1".to_vec()),
930                        Some(b"created".to_vec()),
931                        0,
932                    )],
933                }],
934            }],
935        };
936
937        let bytes = request.encode().unwrap();
938        assert_eq!(
939            &bytes[0..17],
940            &[0, 0, 0, 2, 0, 0, 0, 5, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',]
941        );
942        assert!(bytes.len() > 60);
943    }
944
945    #[test]
946    fn encodes_idempotent_record_batch_identity() {
947        let set = encode_record_batch_set_with_compression_and_identity(
948            &[RecordBatchMessage::new(
949                Some(b"order-1".to_vec()),
950                Some(b"created".to_vec()),
951                1_000,
952            )],
953            RecordBatchCompression::None,
954            RecordBatchIdentity {
955                producer_id: 42,
956                producer_epoch: 3,
957                base_sequence: 7,
958            },
959        )
960        .unwrap();
961
962        assert_eq!(&set[43..51], &42_i64.to_be_bytes());
963        assert_eq!(&set[51..53], &3_i16.to_be_bytes());
964        assert_eq!(&set[53..57], &7_i32.to_be_bytes());
965    }
966
967    #[test]
968    fn encodes_transactional_record_batch_attribute() {
969        let set = encode_record_batch_set_with_compression_identity_and_transaction(
970            &[RecordBatchMessage::new(
971                None,
972                Some(b"created".to_vec()),
973                1_000,
974            )],
975            RecordBatchCompression::None,
976            RecordBatchIdentity {
977                producer_id: 42,
978                producer_epoch: 3,
979                base_sequence: 0,
980            },
981            true,
982        )
983        .unwrap();
984
985        assert_eq!(&set[21..23], &0x10_i16.to_be_bytes());
986    }
987
988    #[test]
989    fn encodes_produce_request_v3_with_record_batch() {
990        let request = ProduceRequestV3 {
991            correlation_id: 5,
992            client_id: Some("kafrust".to_owned()),
993            transactional_id: None,
994            acks: 1,
995            timeout_ms: 30_000,
996            topics: vec![ProduceTopicV3 {
997                name: "orders".to_owned(),
998                partitions: vec![ProducePartitionV3 {
999                    partition_index: 0,
1000                    compression: RecordBatchCompression::None,
1001                    identity: RecordBatchIdentity::NON_IDEMPOTENT,
1002                    records: vec![RecordBatchMessage::new(
1003                        Some(b"order-1".to_vec()),
1004                        Some(b"created".to_vec()),
1005                        1_000,
1006                    )
1007                    .header("source", Some(b"checkout".to_vec()))],
1008                }],
1009            }],
1010        };
1011
1012        let bytes = request.encode().unwrap();
1013
1014        assert_eq!(&bytes[0..4], &[0, 0, 0, 3]);
1015        assert!(bytes.len() > 80);
1016    }
1017
1018    #[test]
1019    fn record_batch_encoding_roundtrips_through_fetch_decoder() {
1020        let record_set = encode_record_batch_set(&[RecordBatchMessage::new(
1021            Some(b"order-1".to_vec()),
1022            Some(b"created".to_vec()),
1023            1_000,
1024        )
1025        .header("source", Some(b"checkout".to_vec()))])
1026        .unwrap();
1027
1028        let mut bytes = Encoder::new();
1029        bytes.write_i32(0);
1030        bytes.write_i32(1);
1031        bytes.write_string("orders").unwrap();
1032        bytes.write_i32(1);
1033        bytes.write_i32(0);
1034        bytes.write_i16(0);
1035        bytes.write_i64(43);
1036        bytes.write_bytes(&record_set).unwrap();
1037        let bytes = bytes.into_bytes();
1038
1039        let mut decoder = Decoder::new(&bytes);
1040        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1041        let record = &response.responses[0].partitions[0].records[0];
1042
1043        assert_eq!(record.offset, 0);
1044        assert_eq!(record.timestamp_ms, 1_000);
1045        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1046        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1047        assert!(decoder.is_empty());
1048    }
1049
1050    #[test]
1051    fn gzip_record_batch_encoding_roundtrips_through_fetch_decoder() {
1052        let record_set =
1053            encode_record_batch_set_with_compression(
1054                &[RecordBatchMessage::new(
1055                    Some(b"order-1".to_vec()),
1056                    Some(b"created".to_vec()),
1057                    1_000,
1058                )
1059                .header("source", Some(b"checkout".to_vec()))],
1060                RecordBatchCompression::Gzip,
1061            )
1062            .unwrap();
1063
1064        let mut bytes = Encoder::new();
1065        bytes.write_i32(0);
1066        bytes.write_i32(1);
1067        bytes.write_string("orders").unwrap();
1068        bytes.write_i32(1);
1069        bytes.write_i32(0);
1070        bytes.write_i16(0);
1071        bytes.write_i64(43);
1072        bytes.write_bytes(&record_set).unwrap();
1073        let bytes = bytes.into_bytes();
1074
1075        let mut decoder = Decoder::new(&bytes);
1076        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1077        let record = &response.responses[0].partitions[0].records[0];
1078
1079        assert_eq!(record.offset, 0);
1080        assert_eq!(record.timestamp_ms, 1_000);
1081        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1082        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1083        assert!(decoder.is_empty());
1084    }
1085
1086    #[test]
1087    fn snappy_record_batch_encoding_roundtrips_through_fetch_decoder() {
1088        let record_set =
1089            encode_record_batch_set_with_compression(
1090                &[RecordBatchMessage::new(
1091                    Some(b"order-1".to_vec()),
1092                    Some(b"created".to_vec()),
1093                    1_000,
1094                )
1095                .header("source", Some(b"checkout".to_vec()))],
1096                RecordBatchCompression::Snappy,
1097            )
1098            .unwrap();
1099
1100        let mut bytes = Encoder::new();
1101        bytes.write_i32(0);
1102        bytes.write_i32(1);
1103        bytes.write_string("orders").unwrap();
1104        bytes.write_i32(1);
1105        bytes.write_i32(0);
1106        bytes.write_i16(0);
1107        bytes.write_i64(43);
1108        bytes.write_bytes(&record_set).unwrap();
1109        let bytes = bytes.into_bytes();
1110
1111        let mut decoder = Decoder::new(&bytes);
1112        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1113        let record = &response.responses[0].partitions[0].records[0];
1114
1115        assert_eq!(record.offset, 0);
1116        assert_eq!(record.timestamp_ms, 1_000);
1117        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1118        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1119        assert!(decoder.is_empty());
1120    }
1121
1122    #[test]
1123    fn lz4_record_batch_encoding_roundtrips_through_fetch_decoder() {
1124        let record_set =
1125            encode_record_batch_set_with_compression(
1126                &[RecordBatchMessage::new(
1127                    Some(b"order-1".to_vec()),
1128                    Some(b"created".to_vec()),
1129                    1_000,
1130                )
1131                .header("source", Some(b"checkout".to_vec()))],
1132                RecordBatchCompression::Lz4,
1133            )
1134            .unwrap();
1135
1136        let mut bytes = Encoder::new();
1137        bytes.write_i32(0);
1138        bytes.write_i32(1);
1139        bytes.write_string("orders").unwrap();
1140        bytes.write_i32(1);
1141        bytes.write_i32(0);
1142        bytes.write_i16(0);
1143        bytes.write_i64(43);
1144        bytes.write_bytes(&record_set).unwrap();
1145        let bytes = bytes.into_bytes();
1146
1147        let mut decoder = Decoder::new(&bytes);
1148        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1149        let record = &response.responses[0].partitions[0].records[0];
1150
1151        assert_eq!(record.offset, 0);
1152        assert_eq!(record.timestamp_ms, 1_000);
1153        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1154        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1155        assert!(decoder.is_empty());
1156    }
1157
1158    #[test]
1159    fn encodes_produce_request_v7_with_record_batch() {
1160        let request = ProduceRequestV7 {
1161            correlation_id: 5,
1162            client_id: Some("kafrust".to_owned()),
1163            transactional_id: None,
1164            acks: 1,
1165            timeout_ms: 30_000,
1166            topics: vec![ProduceTopicV3 {
1167                name: "orders".to_owned(),
1168                partitions: vec![ProducePartitionV3 {
1169                    partition_index: 0,
1170                    compression: RecordBatchCompression::Zstd,
1171                    identity: RecordBatchIdentity::NON_IDEMPOTENT,
1172                    records: vec![RecordBatchMessage::new(
1173                        Some(b"order-1".to_vec()),
1174                        Some(b"created".to_vec()),
1175                        1_000,
1176                    )],
1177                }],
1178            }],
1179        };
1180
1181        let bytes = request.encode().unwrap();
1182
1183        assert_eq!(&bytes[0..4], &[0, 0, 0, 7]);
1184        assert!(bytes.len() > 70);
1185    }
1186
1187    #[test]
1188    fn encodes_produce_request_v9_with_flexible_record_batch() {
1189        let request = ProduceRequestV9 {
1190            correlation_id: 5,
1191            client_id: Some("kafrust".to_owned()),
1192            transactional_id: Some("orders-tx".to_owned()),
1193            acks: -1,
1194            timeout_ms: 30_000,
1195            topics: vec![ProduceTopicV3 {
1196                name: "orders".to_owned(),
1197                partitions: vec![ProducePartitionV3 {
1198                    partition_index: 0,
1199                    compression: RecordBatchCompression::None,
1200                    identity: RecordBatchIdentity {
1201                        producer_id: 42,
1202                        producer_epoch: 3,
1203                        base_sequence: 7,
1204                    },
1205                    records: vec![RecordBatchMessage::new(
1206                        Some(b"order-1".to_vec()),
1207                        Some(b"created".to_vec()),
1208                        1_000,
1209                    )],
1210                }],
1211            }],
1212        };
1213
1214        let bytes = request.encode().unwrap();
1215
1216        assert_eq!(&bytes[0..4], &[0, 0, 0, 9]);
1217        assert!(bytes
1218            .windows(7)
1219            .any(|window| window == [10, b'o', b'r', b'd', b'e', b'r', b's']));
1220        assert!(bytes.len() > 80);
1221    }
1222
1223    #[test]
1224    fn encodes_produce_request_v11_with_flexible_record_batch_schema() {
1225        let request = ProduceRequestV11 {
1226            correlation_id: 5,
1227            client_id: Some("kafrust".to_owned()),
1228            transactional_id: None,
1229            acks: 1,
1230            timeout_ms: 30_000,
1231            topics: Vec::new(),
1232        };
1233
1234        let bytes = request.encode().unwrap();
1235
1236        assert_eq!(&bytes[0..4], &[0, 0, 0, 11]);
1237        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1238        assert!(bytes.len() > 20);
1239    }
1240
1241    #[test]
1242    fn encodes_produce_request_v12_with_flexible_record_batch_schema() {
1243        let request = ProduceRequestV12 {
1244            correlation_id: 5,
1245            client_id: Some("kafrust".to_owned()),
1246            transactional_id: None,
1247            acks: 1,
1248            timeout_ms: 30_000,
1249            topics: Vec::new(),
1250        };
1251
1252        let bytes = request.encode().unwrap();
1253
1254        assert_eq!(&bytes[0..4], &[0, 0, 0, 12]);
1255        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1256        assert!(bytes.len() > 20);
1257    }
1258
1259    #[test]
1260    fn encodes_produce_request_v13_with_topic_uuid() {
1261        let request = ProduceRequestV13 {
1262            correlation_id: 5,
1263            client_id: Some("kafrust".to_owned()),
1264            transactional_id: None,
1265            acks: 1,
1266            timeout_ms: 30_000,
1267            topics: vec![ProduceTopicV13 {
1268                topic_id: [7; 16],
1269                partitions: Vec::new(),
1270            }],
1271        };
1272
1273        let bytes = request.encode().unwrap();
1274
1275        assert_eq!(&bytes[0..4], &[0, 0, 0, 13]);
1276        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1277        assert!(bytes.windows(16).any(|window| window == [7; 16]));
1278        assert!(bytes.len() > 35);
1279    }
1280
1281    #[test]
1282    fn zstd_record_batch_encoding_roundtrips_through_fetch_decoder() {
1283        let record_set =
1284            encode_record_batch_set_with_compression(
1285                &[RecordBatchMessage::new(
1286                    Some(b"order-1".to_vec()),
1287                    Some(b"created".to_vec()),
1288                    1_000,
1289                )
1290                .header("source", Some(b"checkout".to_vec()))],
1291                RecordBatchCompression::Zstd,
1292            )
1293            .unwrap();
1294
1295        let mut bytes = Encoder::new();
1296        bytes.write_i32(0);
1297        bytes.write_i32(1);
1298        bytes.write_string("orders").unwrap();
1299        bytes.write_i32(1);
1300        bytes.write_i32(0);
1301        bytes.write_i16(0);
1302        bytes.write_i64(43);
1303        bytes.write_bytes(&record_set).unwrap();
1304        let bytes = bytes.into_bytes();
1305
1306        let mut decoder = Decoder::new(&bytes);
1307        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1308        let record = &response.responses[0].partitions[0].records[0];
1309
1310        assert_eq!(record.offset, 0);
1311        assert_eq!(record.timestamp_ms, 1_000);
1312        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1313        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1314        assert!(decoder.is_empty());
1315    }
1316
1317    #[test]
1318    fn fetch_decoder_applies_custom_decompression_limit_to_record_batch() {
1319        let record_set = encode_record_batch_set_with_compression(
1320            &[RecordBatchMessage::new(
1321                Some(b"order-1".to_vec()),
1322                Some(vec![b'x'; 1024]),
1323                1_000,
1324            )],
1325            RecordBatchCompression::Zstd,
1326        )
1327        .unwrap();
1328
1329        let mut bytes = Encoder::new();
1330        bytes.write_i32(0);
1331        bytes.write_i32(1);
1332        bytes.write_string("orders").unwrap();
1333        bytes.write_i32(1);
1334        bytes.write_i32(0);
1335        bytes.write_i16(0);
1336        bytes.write_i64(43);
1337        bytes.write_bytes(&record_set).unwrap();
1338        let bytes = bytes.into_bytes();
1339        let limits = DecodeLimits::new().with_max_decompressed_record_bytes(64);
1340        let mut decoder = Decoder::with_limits(&bytes, limits);
1341
1342        assert!(matches!(
1343            FetchResponseV2::decode_body(&mut decoder),
1344            Err(crate::Error::LimitExceeded {
1345                kind: "decompressed record batch bytes",
1346                max: 64,
1347                ..
1348            })
1349        ));
1350    }
1351
1352    #[test]
1353    fn reports_message_set_encoded_len() {
1354        let records = [MessageSetMessage::new(
1355            Some(b"order-1".to_vec()),
1356            Some(b"created".to_vec()),
1357            1_000,
1358        )];
1359
1360        assert_eq!(
1361            encoded_message_set_len(&records).unwrap(),
1362            encode_message_set(&records).unwrap().len()
1363        );
1364    }
1365
1366    #[test]
1367    fn reports_record_batch_set_encoded_len() {
1368        let records =
1369            [
1370                RecordBatchMessage::new(
1371                    Some(b"order-1".to_vec()),
1372                    Some(b"created".to_vec()),
1373                    1_000,
1374                )
1375                .header("source", Some(b"checkout".to_vec())),
1376            ];
1377
1378        assert_eq!(
1379            encoded_record_batch_set_len(&records).unwrap(),
1380            encode_record_batch_set(&records).unwrap().len()
1381        );
1382    }
1383
1384    #[test]
1385    fn decodes_produce_response_v2() {
1386        let bytes = [
1387            0, 0, 0, 1, // topic response count
1388            0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
1389            0, 0, 0, 1, // partition response count
1390            0, 0, 0, 0, // partition
1391            0, 0, // error code
1392            0, 0, 0, 0, 0, 0, 0, 42, // base offset
1393            0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, // log append time -1
1394            0, 0, 0, 0, // throttle time
1395        ];
1396        let mut decoder = Decoder::new(&bytes);
1397        let response = ProduceResponseV2::decode_body(&mut decoder).unwrap();
1398
1399        assert_eq!(response.throttle_time_ms, 0);
1400        assert_eq!(response.responses[0].name, "orders");
1401        assert_eq!(response.responses[0].partitions[0].partition_index, 0);
1402        assert_eq!(response.responses[0].partitions[0].error_code, 0);
1403        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1404        assert!(decoder.is_empty());
1405    }
1406
1407    #[test]
1408    fn decodes_produce_response_v7() {
1409        let bytes = [
1410            0, 0, 0, 1, // topic response count
1411            0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
1412            0, 0, 0, 1, // partition response count
1413            0, 0, 0, 0, // partition
1414            0, 0, // error code
1415            0, 0, 0, 0, 0, 0, 0, 42, // base offset
1416            0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, // log append time -1
1417            0, 0, 0, 0, 0, 0, 0, 7, // log start offset
1418            0, 0, 0, 0, // throttle time
1419        ];
1420        let mut decoder = Decoder::new(&bytes);
1421        let response = ProduceResponseV7::decode_body(&mut decoder).unwrap();
1422
1423        assert_eq!(response.throttle_time_ms, 0);
1424        assert_eq!(response.responses[0].name, "orders");
1425        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1426        assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1427        assert!(decoder.is_empty());
1428    }
1429
1430    #[test]
1431    fn decodes_produce_response_v9_with_record_error_fields() {
1432        let mut bytes = Encoder::new();
1433        bytes
1434            .write_compact_array(Some(&[()]), |encoder, _| {
1435                encoder.write_compact_string("orders")?;
1436                encoder.write_compact_array(Some(&[()]), |encoder, _| {
1437                    encoder.write_i32(0);
1438                    encoder.write_i16(0);
1439                    encoder.write_i64(42);
1440                    encoder.write_i64(-1);
1441                    encoder.write_i64(7);
1442                    encoder.write_compact_array(Some(&[()]), |encoder, _| {
1443                        encoder.write_i32(3);
1444                        encoder.write_compact_nullable_string(Some("bad record"))?;
1445                        encoder.write_empty_tagged_fields();
1446                        Ok(())
1447                    })?;
1448                    encoder.write_compact_nullable_string(Some("batch rejected"))?;
1449                    encoder.write_empty_tagged_fields();
1450                    Ok(())
1451                })?;
1452                encoder.write_empty_tagged_fields();
1453                Ok(())
1454            })
1455            .unwrap();
1456        bytes.write_i32(0);
1457        bytes.write_empty_tagged_fields();
1458
1459        let bytes = bytes.into_bytes();
1460        let mut decoder = Decoder::new(&bytes);
1461        let response = ProduceResponseV9::decode_body(&mut decoder).unwrap();
1462
1463        assert_eq!(response.responses[0].name, "orders");
1464        let partition = &response.responses[0].partitions[0];
1465        assert_eq!(partition.base_offset, 42);
1466        assert_eq!(partition.log_start_offset, 7);
1467        assert_eq!(partition.record_errors[0].batch_index, 3);
1468        assert_eq!(
1469            partition.record_errors[0].error_message.as_deref(),
1470            Some("bad record")
1471        );
1472        assert_eq!(partition.error_message.as_deref(), Some("batch rejected"));
1473        assert!(decoder.is_empty());
1474    }
1475
1476    #[test]
1477    fn decodes_produce_response_v13_with_topic_uuid() {
1478        let mut bytes = Encoder::new();
1479        bytes
1480            .write_compact_array(Some(&[()]), |encoder, _| {
1481                encoder.write_uuid(&[7; 16]);
1482                encoder.write_compact_array(Some(&[()]), |encoder, _| {
1483                    encoder.write_i32(0);
1484                    encoder.write_i16(0);
1485                    encoder.write_i64(42);
1486                    encoder.write_i64(-1);
1487                    encoder.write_i64(7);
1488                    encoder.write_unsigned_varint(1);
1489                    encoder.write_compact_nullable_string(None)?;
1490                    encoder.write_empty_tagged_fields();
1491                    Ok(())
1492                })?;
1493                encoder.write_empty_tagged_fields();
1494                Ok(())
1495            })
1496            .unwrap();
1497        bytes.write_i32(0);
1498        bytes.write_empty_tagged_fields();
1499
1500        let bytes = bytes.into_bytes();
1501        let mut decoder = Decoder::new(&bytes);
1502        let response = ProduceResponseV13::decode_body(&mut decoder).unwrap();
1503
1504        assert_eq!(response.responses[0].topic_id, [7; 16]);
1505        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1506        assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1507        assert!(decoder.is_empty());
1508    }
1509}