Skip to main content

kafrust_protocol/api/
produce.rs

1use crate::codec::{DecodeLimits, 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    /// Broker endpoints advertised by the v10+ tagged response field.
528    pub node_endpoints: Vec<ProduceNodeEndpointV10>,
529}
530
531/// Flexible Produce response used by Kafka API versions 11 and newer.
532pub type ProduceResponseV11 = ProduceResponseV9;
533
534/// Flexible Produce response for Kafka API version 12.
535pub type ProduceResponseV12 = ProduceResponseV9;
536
537/// Flexible Produce response for Kafka API version 13.
538#[derive(Debug, Clone, PartialEq, Eq)]
539pub struct ProduceResponseV13 {
540    pub responses: Vec<ProduceTopicResponseV13>,
541    pub throttle_time_ms: i32,
542    /// Broker endpoints advertised by the v10+ tagged response field.
543    pub node_endpoints: Vec<ProduceNodeEndpointV10>,
544}
545
546impl ProduceResponseV13 {
547    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
548        let responses = decoder
549            .read_compact_array("produce responses", ProduceTopicResponseV13::decode)?
550            .unwrap_or_default();
551        let throttle_time_ms = decoder.read_i32()?;
552        let tagged_fields = decoder.read_tagged_fields()?;
553        let node_endpoints = tagged_field_data(&tagged_fields, 0)
554            .map(|data| decode_produce_node_endpoints(data, decoder.limits()))
555            .transpose()?
556            .unwrap_or_default();
557        Ok(Self {
558            responses,
559            throttle_time_ms,
560            node_endpoints,
561        })
562    }
563}
564
565#[derive(Debug, Clone, PartialEq, Eq)]
566pub struct ProduceTopicResponseV13 {
567    pub topic_id: [u8; 16],
568    pub partitions: Vec<ProducePartitionResponseV9>,
569}
570
571/// The current partition leader advertised by Produce response v10 and newer.
572#[derive(Debug, Clone, Copy, PartialEq, Eq)]
573pub struct ProduceLeaderIdAndEpochV10 {
574    pub leader_id: i32,
575    pub leader_epoch: i32,
576}
577
578/// A broker endpoint advertised by Produce response v10 and newer.
579#[derive(Debug, Clone, PartialEq, Eq)]
580pub struct ProduceNodeEndpointV10 {
581    pub node_id: i32,
582    pub host: String,
583    pub port: i32,
584    pub rack: Option<String>,
585}
586
587impl ProduceTopicResponseV13 {
588    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
589        let topic_id = decoder.read_uuid()?;
590        let partitions = decoder
591            .read_compact_array(
592                "produce partition responses",
593                ProducePartitionResponseV9::decode,
594            )?
595            .unwrap_or_default();
596        decoder.read_tagged_fields()?;
597        Ok(Self {
598            topic_id,
599            partitions,
600        })
601    }
602}
603
604impl ProduceResponseV9 {
605    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
606        let responses = decoder
607            .read_compact_array("produce responses", ProduceTopicResponseV9::decode)?
608            .unwrap_or_default();
609        let throttle_time_ms = decoder.read_i32()?;
610        let tagged_fields = decoder.read_tagged_fields()?;
611        let node_endpoints = tagged_field_data(&tagged_fields, 0)
612            .map(|data| decode_produce_node_endpoints(data, decoder.limits()))
613            .transpose()?
614            .unwrap_or_default();
615        Ok(Self {
616            responses,
617            throttle_time_ms,
618            node_endpoints,
619        })
620    }
621}
622
623fn tagged_field_data(fields: &[crate::codec::TaggedField], tag: u32) -> Option<&[u8]> {
624    fields
625        .iter()
626        .find(|field| field.tag == tag)
627        .map(|field| field.data.as_slice())
628}
629
630fn decode_produce_node_endpoints(
631    data: &[u8],
632    limits: DecodeLimits,
633) -> Result<Vec<ProduceNodeEndpointV10>> {
634    let mut decoder = Decoder::with_limits(data, limits);
635    Ok(decoder
636        .read_compact_array("produce node endpoints", |decoder| {
637            let node_id = decoder.read_i32()?;
638            let host = decoder.read_compact_string()?;
639            let port = decoder.read_i32()?;
640            let rack = decoder.read_compact_nullable_string()?;
641            decoder.read_tagged_fields()?;
642            Ok(ProduceNodeEndpointV10 {
643                node_id,
644                host,
645                port,
646                rack,
647            })
648        })?
649        .unwrap_or_default())
650}
651
652#[derive(Debug, Clone, PartialEq, Eq)]
653pub struct ProduceTopicResponseV9 {
654    pub name: String,
655    pub partitions: Vec<ProducePartitionResponseV9>,
656}
657
658impl ProduceTopicResponseV9 {
659    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
660        let name = decoder.read_compact_string()?;
661        let partitions = decoder
662            .read_compact_array(
663                "produce partition responses",
664                ProducePartitionResponseV9::decode,
665            )?
666            .unwrap_or_default();
667        decoder.read_tagged_fields()?;
668        Ok(Self { name, partitions })
669    }
670}
671
672#[derive(Debug, Clone, PartialEq, Eq)]
673pub struct ProducePartitionResponseV9 {
674    pub partition_index: i32,
675    pub error_code: i16,
676    pub base_offset: i64,
677    pub log_append_time_ms: i64,
678    pub log_start_offset: i64,
679    pub record_errors: Vec<ProduceRecordErrorV9>,
680    pub error_message: Option<String>,
681    /// The broker's current leader hint, when present in the v10+ tag set.
682    pub current_leader: Option<ProduceLeaderIdAndEpochV10>,
683}
684
685impl ProducePartitionResponseV9 {
686    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
687        let partition_index = decoder.read_i32()?;
688        let error_code = decoder.read_i16()?;
689        let base_offset = decoder.read_i64()?;
690        let log_append_time_ms = decoder.read_i64()?;
691        let log_start_offset = decoder.read_i64()?;
692        let record_errors = decoder
693            .read_compact_array("produce record errors", ProduceRecordErrorV9::decode)?
694            .unwrap_or_default();
695        let error_message = decoder.read_compact_nullable_string()?;
696        let tagged_fields = decoder.read_tagged_fields()?;
697        let current_leader = tagged_field_data(&tagged_fields, 0)
698            .map(|data| decode_produce_current_leader(data, decoder.limits()))
699            .transpose()?;
700        Ok(Self {
701            partition_index,
702            error_code,
703            base_offset,
704            log_append_time_ms,
705            log_start_offset,
706            record_errors,
707            error_message,
708            current_leader,
709        })
710    }
711}
712
713fn decode_produce_current_leader(
714    data: &[u8],
715    limits: DecodeLimits,
716) -> Result<ProduceLeaderIdAndEpochV10> {
717    let mut decoder = Decoder::with_limits(data, limits);
718    let leader_id = decoder.read_i32()?;
719    let leader_epoch = decoder.read_i32()?;
720    decoder.read_tagged_fields()?;
721    Ok(ProduceLeaderIdAndEpochV10 {
722        leader_id,
723        leader_epoch,
724    })
725}
726
727#[derive(Debug, Clone, PartialEq, Eq)]
728pub struct ProduceRecordErrorV9 {
729    pub batch_index: i32,
730    pub error_message: Option<String>,
731}
732
733impl ProduceRecordErrorV9 {
734    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
735        let batch_index = decoder.read_i32()?;
736        let error_message = decoder.read_compact_nullable_string()?;
737        decoder.read_tagged_fields()?;
738        Ok(Self {
739            batch_index,
740            error_message,
741        })
742    }
743}
744
745impl ProduceResponseV7 {
746    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
747        Ok(Self {
748            responses: decoder
749                .read_array("produce responses", ProduceTopicResponseV7::decode)?
750                .unwrap_or_default(),
751            throttle_time_ms: decoder.read_i32()?,
752        })
753    }
754}
755
756#[derive(Debug, Clone, PartialEq, Eq)]
757pub struct ProduceTopicResponseV7 {
758    pub name: String,
759    pub partitions: Vec<ProducePartitionResponseV7>,
760}
761
762impl ProduceTopicResponseV7 {
763    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
764        Ok(Self {
765            name: decoder.read_string()?,
766            partitions: decoder
767                .read_array(
768                    "produce partition responses",
769                    ProducePartitionResponseV7::decode,
770                )?
771                .unwrap_or_default(),
772        })
773    }
774}
775
776#[derive(Debug, Clone, PartialEq, Eq)]
777pub struct ProducePartitionResponseV7 {
778    pub partition_index: i32,
779    pub error_code: i16,
780    pub base_offset: i64,
781    pub log_append_time_ms: i64,
782    pub log_start_offset: i64,
783}
784
785impl ProducePartitionResponseV7 {
786    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
787        Ok(Self {
788            partition_index: decoder.read_i32()?,
789            error_code: decoder.read_i16()?,
790            base_offset: decoder.read_i64()?,
791            log_append_time_ms: decoder.read_i64()?,
792            log_start_offset: decoder.read_i64()?,
793        })
794    }
795}
796
797fn encode_message_set(records: &[MessageSetMessage]) -> Result<Vec<u8>> {
798    let mut set = Encoder::new();
799    for record in records {
800        let message = encode_message(record)?;
801        set.write_i64(0);
802        set.write_i32(i32::try_from(message.len()).map_err(|_| Error::LengthOverflow("message"))?);
803        set.write_raw(&message);
804    }
805    Ok(set.into_bytes())
806}
807
808fn encode_message(record: &MessageSetMessage) -> Result<Vec<u8>> {
809    let mut body = Encoder::new();
810    body.write_i8(1);
811    body.write_i8(0);
812    body.write_i64(record.timestamp_ms);
813    body.write_nullable_bytes(record.key.as_deref())?;
814    body.write_nullable_bytes(record.value.as_deref())?;
815    let body = body.into_bytes();
816
817    let mut message = Encoder::new();
818    message.write_i32(crc32_ieee(&body) as i32);
819    message.write_raw(&body);
820    Ok(message.into_bytes())
821}
822
823fn encode_record_batch_set(records: &[RecordBatchMessage]) -> Result<Vec<u8>> {
824    encode_record_batch_set_with_compression(records, RecordBatchCompression::None)
825}
826
827fn encode_record_batch_set_with_compression(
828    records: &[RecordBatchMessage],
829    compression: RecordBatchCompression,
830) -> Result<Vec<u8>> {
831    encode_record_batch_set_with_compression_and_identity(
832        records,
833        compression,
834        RecordBatchIdentity::NON_IDEMPOTENT,
835    )
836}
837
838fn encode_record_batch_set_with_compression_and_identity(
839    records: &[RecordBatchMessage],
840    compression: RecordBatchCompression,
841    identity: RecordBatchIdentity,
842) -> Result<Vec<u8>> {
843    encode_record_batch_set_with_compression_identity_and_transaction(
844        records,
845        compression,
846        identity,
847        false,
848    )
849}
850
851fn encode_record_batch_set_with_compression_identity_and_transaction(
852    records: &[RecordBatchMessage],
853    compression: RecordBatchCompression,
854    identity: RecordBatchIdentity,
855    transactional: bool,
856) -> Result<Vec<u8>> {
857    let base_timestamp = records
858        .first()
859        .map(|record| record.timestamp_ms)
860        .unwrap_or_default();
861    let max_timestamp = records
862        .iter()
863        .map(|record| record.timestamp_ms)
864        .max()
865        .unwrap_or(base_timestamp);
866    let last_offset_delta = records
867        .len()
868        .checked_sub(1)
869        .map(|delta| i32::try_from(delta).map_err(|_| Error::LengthOverflow("record batch")))
870        .transpose()?
871        .unwrap_or_default();
872
873    let record_count =
874        i32::try_from(records.len()).map_err(|_| Error::LengthOverflow("record batch records"))?;
875    let mut record_bytes = Encoder::new();
876    for (offset_delta, record) in records.iter().enumerate() {
877        let encoded = encode_record(record, base_timestamp, offset_delta)?;
878        record_bytes.write_varint(
879            i32::try_from(encoded.len()).map_err(|_| Error::LengthOverflow("record"))?,
880        );
881        record_bytes.write_raw(&encoded);
882    }
883    let record_bytes = compress_record_batch_records(compression, &record_bytes.into_bytes())?;
884
885    let mut crc_payload = Encoder::new();
886    let attributes = compression.attributes() | if transactional { 0x10 } else { 0 };
887    crc_payload.write_i16(attributes);
888    crc_payload.write_i32(last_offset_delta);
889    crc_payload.write_i64(base_timestamp);
890    crc_payload.write_i64(max_timestamp);
891    crc_payload.write_i64(identity.producer_id);
892    crc_payload.write_i16(identity.producer_epoch);
893    crc_payload.write_i32(identity.base_sequence);
894    crc_payload.write_i32(record_count);
895    crc_payload.write_raw(&record_bytes);
896    let crc_payload = crc_payload.into_bytes();
897
898    let mut batch = Encoder::new();
899    batch.write_i32(0);
900    batch.write_i8(2);
901    batch.write_i32(crc32c(&crc_payload) as i32);
902    batch.write_raw(&crc_payload);
903    let batch = batch.into_bytes();
904
905    let mut set = Encoder::new();
906    set.write_i64(0);
907    set.write_i32(i32::try_from(batch.len()).map_err(|_| Error::LengthOverflow("record batch"))?);
908    set.write_raw(&batch);
909    Ok(set.into_bytes())
910}
911
912fn encode_record(
913    record: &RecordBatchMessage,
914    base_timestamp: i64,
915    offset_delta: usize,
916) -> Result<Vec<u8>> {
917    let mut encoder = Encoder::new();
918    encoder.write_i8(0);
919    encoder.write_varlong(record.timestamp_ms.saturating_sub(base_timestamp));
920    encoder.write_varint(
921        i32::try_from(offset_delta).map_err(|_| Error::LengthOverflow("record offset delta"))?,
922    );
923    encoder.write_varint_nullable_bytes(record.key.as_deref())?;
924    encoder.write_varint_nullable_bytes(record.value.as_deref())?;
925    encoder.write_varint(
926        i32::try_from(record.headers.len()).map_err(|_| Error::LengthOverflow("record headers"))?,
927    );
928    for header in &record.headers {
929        encoder.write_varint_bytes(header.key.as_bytes())?;
930        encoder.write_varint_nullable_bytes(header.value.as_deref())?;
931    }
932    Ok(encoder.into_bytes())
933}
934
935fn crc32_ieee(bytes: &[u8]) -> u32 {
936    crc32_with_table(bytes, &CRC32_IEEE_TABLE)
937}
938
939fn crc32c(bytes: &[u8]) -> u32 {
940    crc32_with_table(bytes, &CRC32C_TABLE)
941}
942
943const CRC32_IEEE_TABLE: [u32; 256] = crc32_table(0xedb8_8320);
944const CRC32C_TABLE: [u32; 256] = crc32_table(0x82f6_3b78);
945
946const fn crc32_table(polynomial: u32) -> [u32; 256] {
947    let mut table = [0; 256];
948    let mut index = 0;
949    while index < table.len() {
950        let mut value = index as u32;
951        let mut bit = 0;
952        while bit < 8 {
953            let mask = 0u32.wrapping_sub(value & 1);
954            value = (value >> 1) ^ (polynomial & mask);
955            bit += 1;
956        }
957        table[index] = value;
958        index += 1;
959    }
960    table
961}
962
963fn crc32_with_table(bytes: &[u8], table: &[u32; 256]) -> u32 {
964    let mut crc = 0xffff_ffffu32;
965    for byte in bytes {
966        let index = usize::from((crc as u8) ^ byte);
967        crc = (crc >> 8) ^ table[index];
968    }
969    !crc
970}
971
972#[cfg(test)]
973#[allow(clippy::unwrap_used)]
974mod tests {
975    use super::{
976        crc32_ieee, crc32c, encode_message_set, encode_record_batch_set,
977        encode_record_batch_set_with_compression,
978        encode_record_batch_set_with_compression_and_identity,
979        encode_record_batch_set_with_compression_identity_and_transaction, encoded_message_set_len,
980        encoded_record_batch_set_len, MessageSetMessage, ProduceLeaderIdAndEpochV10,
981        ProducePartitionV2, ProducePartitionV3, ProduceRequestV11, ProduceRequestV12,
982        ProduceRequestV13, ProduceRequestV2, ProduceRequestV3, ProduceRequestV7, ProduceRequestV9,
983        ProduceResponseV13, ProduceResponseV2, ProduceResponseV7, ProduceResponseV9,
984        ProduceTopicV13, ProduceTopicV2, ProduceTopicV3, RecordBatchIdentity, RecordBatchMessage,
985    };
986    use crate::codec::{DecodeLimits, Decoder};
987    use crate::record_batch::RecordBatchCompression;
988    use crate::{api::fetch::FetchResponseV2, codec::Encoder};
989
990    #[test]
991    fn crc_implementations_match_standard_check_vectors() {
992        assert_eq!(crc32_ieee(b"123456789"), 0xcbf4_3926);
993        assert_eq!(crc32c(b"123456789"), 0xe306_9283);
994    }
995
996    #[test]
997    fn encodes_produce_request_v2() {
998        let request = ProduceRequestV2 {
999            correlation_id: 5,
1000            client_id: Some("kafrust".to_owned()),
1001            acks: 1,
1002            timeout_ms: 30_000,
1003            topics: vec![ProduceTopicV2 {
1004                name: "orders".to_owned(),
1005                partitions: vec![ProducePartitionV2 {
1006                    partition_index: 0,
1007                    records: vec![MessageSetMessage::new(
1008                        Some(b"order-1".to_vec()),
1009                        Some(b"created".to_vec()),
1010                        0,
1011                    )],
1012                }],
1013            }],
1014        };
1015
1016        let bytes = request.encode().unwrap();
1017        assert_eq!(
1018            &bytes[0..17],
1019            &[0, 0, 0, 2, 0, 0, 0, 5, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',]
1020        );
1021        assert!(bytes.len() > 60);
1022    }
1023
1024    #[test]
1025    fn encodes_idempotent_record_batch_identity() {
1026        let set = encode_record_batch_set_with_compression_and_identity(
1027            &[RecordBatchMessage::new(
1028                Some(b"order-1".to_vec()),
1029                Some(b"created".to_vec()),
1030                1_000,
1031            )],
1032            RecordBatchCompression::None,
1033            RecordBatchIdentity {
1034                producer_id: 42,
1035                producer_epoch: 3,
1036                base_sequence: 7,
1037            },
1038        )
1039        .unwrap();
1040
1041        assert_eq!(&set[43..51], &42_i64.to_be_bytes());
1042        assert_eq!(&set[51..53], &3_i16.to_be_bytes());
1043        assert_eq!(&set[53..57], &7_i32.to_be_bytes());
1044    }
1045
1046    #[test]
1047    fn encodes_transactional_record_batch_attribute() {
1048        let set = encode_record_batch_set_with_compression_identity_and_transaction(
1049            &[RecordBatchMessage::new(
1050                None,
1051                Some(b"created".to_vec()),
1052                1_000,
1053            )],
1054            RecordBatchCompression::None,
1055            RecordBatchIdentity {
1056                producer_id: 42,
1057                producer_epoch: 3,
1058                base_sequence: 0,
1059            },
1060            true,
1061        )
1062        .unwrap();
1063
1064        assert_eq!(&set[21..23], &0x10_i16.to_be_bytes());
1065    }
1066
1067    #[test]
1068    fn encodes_produce_request_v3_with_record_batch() {
1069        let request = ProduceRequestV3 {
1070            correlation_id: 5,
1071            client_id: Some("kafrust".to_owned()),
1072            transactional_id: None,
1073            acks: 1,
1074            timeout_ms: 30_000,
1075            topics: vec![ProduceTopicV3 {
1076                name: "orders".to_owned(),
1077                partitions: vec![ProducePartitionV3 {
1078                    partition_index: 0,
1079                    compression: RecordBatchCompression::None,
1080                    identity: RecordBatchIdentity::NON_IDEMPOTENT,
1081                    records: vec![RecordBatchMessage::new(
1082                        Some(b"order-1".to_vec()),
1083                        Some(b"created".to_vec()),
1084                        1_000,
1085                    )
1086                    .header("source", Some(b"checkout".to_vec()))],
1087                }],
1088            }],
1089        };
1090
1091        let bytes = request.encode().unwrap();
1092
1093        assert_eq!(&bytes[0..4], &[0, 0, 0, 3]);
1094        assert!(bytes.len() > 80);
1095    }
1096
1097    #[test]
1098    fn record_batch_encoding_roundtrips_through_fetch_decoder() {
1099        let record_set = encode_record_batch_set(&[RecordBatchMessage::new(
1100            Some(b"order-1".to_vec()),
1101            Some(b"created".to_vec()),
1102            1_000,
1103        )
1104        .header("source", Some(b"checkout".to_vec()))])
1105        .unwrap();
1106
1107        let mut bytes = Encoder::new();
1108        bytes.write_i32(0);
1109        bytes.write_i32(1);
1110        bytes.write_string("orders").unwrap();
1111        bytes.write_i32(1);
1112        bytes.write_i32(0);
1113        bytes.write_i16(0);
1114        bytes.write_i64(43);
1115        bytes.write_bytes(&record_set).unwrap();
1116        let bytes = bytes.into_bytes();
1117
1118        let mut decoder = Decoder::new(&bytes);
1119        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1120        let record = &response.responses[0].partitions[0].records[0];
1121
1122        assert_eq!(record.offset, 0);
1123        assert_eq!(record.timestamp_ms, 1_000);
1124        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1125        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1126        assert!(decoder.is_empty());
1127    }
1128
1129    #[test]
1130    fn gzip_record_batch_encoding_roundtrips_through_fetch_decoder() {
1131        let record_set =
1132            encode_record_batch_set_with_compression(
1133                &[RecordBatchMessage::new(
1134                    Some(b"order-1".to_vec()),
1135                    Some(b"created".to_vec()),
1136                    1_000,
1137                )
1138                .header("source", Some(b"checkout".to_vec()))],
1139                RecordBatchCompression::Gzip,
1140            )
1141            .unwrap();
1142
1143        let mut bytes = Encoder::new();
1144        bytes.write_i32(0);
1145        bytes.write_i32(1);
1146        bytes.write_string("orders").unwrap();
1147        bytes.write_i32(1);
1148        bytes.write_i32(0);
1149        bytes.write_i16(0);
1150        bytes.write_i64(43);
1151        bytes.write_bytes(&record_set).unwrap();
1152        let bytes = bytes.into_bytes();
1153
1154        let mut decoder = Decoder::new(&bytes);
1155        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1156        let record = &response.responses[0].partitions[0].records[0];
1157
1158        assert_eq!(record.offset, 0);
1159        assert_eq!(record.timestamp_ms, 1_000);
1160        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1161        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1162        assert!(decoder.is_empty());
1163    }
1164
1165    #[test]
1166    fn snappy_record_batch_encoding_roundtrips_through_fetch_decoder() {
1167        let record_set =
1168            encode_record_batch_set_with_compression(
1169                &[RecordBatchMessage::new(
1170                    Some(b"order-1".to_vec()),
1171                    Some(b"created".to_vec()),
1172                    1_000,
1173                )
1174                .header("source", Some(b"checkout".to_vec()))],
1175                RecordBatchCompression::Snappy,
1176            )
1177            .unwrap();
1178
1179        let mut bytes = Encoder::new();
1180        bytes.write_i32(0);
1181        bytes.write_i32(1);
1182        bytes.write_string("orders").unwrap();
1183        bytes.write_i32(1);
1184        bytes.write_i32(0);
1185        bytes.write_i16(0);
1186        bytes.write_i64(43);
1187        bytes.write_bytes(&record_set).unwrap();
1188        let bytes = bytes.into_bytes();
1189
1190        let mut decoder = Decoder::new(&bytes);
1191        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1192        let record = &response.responses[0].partitions[0].records[0];
1193
1194        assert_eq!(record.offset, 0);
1195        assert_eq!(record.timestamp_ms, 1_000);
1196        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1197        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1198        assert!(decoder.is_empty());
1199    }
1200
1201    #[test]
1202    fn lz4_record_batch_encoding_roundtrips_through_fetch_decoder() {
1203        let record_set =
1204            encode_record_batch_set_with_compression(
1205                &[RecordBatchMessage::new(
1206                    Some(b"order-1".to_vec()),
1207                    Some(b"created".to_vec()),
1208                    1_000,
1209                )
1210                .header("source", Some(b"checkout".to_vec()))],
1211                RecordBatchCompression::Lz4,
1212            )
1213            .unwrap();
1214
1215        let mut bytes = Encoder::new();
1216        bytes.write_i32(0);
1217        bytes.write_i32(1);
1218        bytes.write_string("orders").unwrap();
1219        bytes.write_i32(1);
1220        bytes.write_i32(0);
1221        bytes.write_i16(0);
1222        bytes.write_i64(43);
1223        bytes.write_bytes(&record_set).unwrap();
1224        let bytes = bytes.into_bytes();
1225
1226        let mut decoder = Decoder::new(&bytes);
1227        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1228        let record = &response.responses[0].partitions[0].records[0];
1229
1230        assert_eq!(record.offset, 0);
1231        assert_eq!(record.timestamp_ms, 1_000);
1232        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1233        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1234        assert!(decoder.is_empty());
1235    }
1236
1237    #[test]
1238    fn encodes_produce_request_v7_with_record_batch() {
1239        let request = ProduceRequestV7 {
1240            correlation_id: 5,
1241            client_id: Some("kafrust".to_owned()),
1242            transactional_id: None,
1243            acks: 1,
1244            timeout_ms: 30_000,
1245            topics: vec![ProduceTopicV3 {
1246                name: "orders".to_owned(),
1247                partitions: vec![ProducePartitionV3 {
1248                    partition_index: 0,
1249                    compression: RecordBatchCompression::Zstd,
1250                    identity: RecordBatchIdentity::NON_IDEMPOTENT,
1251                    records: vec![RecordBatchMessage::new(
1252                        Some(b"order-1".to_vec()),
1253                        Some(b"created".to_vec()),
1254                        1_000,
1255                    )],
1256                }],
1257            }],
1258        };
1259
1260        let bytes = request.encode().unwrap();
1261
1262        assert_eq!(&bytes[0..4], &[0, 0, 0, 7]);
1263        assert!(bytes.len() > 70);
1264    }
1265
1266    #[test]
1267    fn encodes_produce_request_v9_with_flexible_record_batch() {
1268        let request = ProduceRequestV9 {
1269            correlation_id: 5,
1270            client_id: Some("kafrust".to_owned()),
1271            transactional_id: Some("orders-tx".to_owned()),
1272            acks: -1,
1273            timeout_ms: 30_000,
1274            topics: vec![ProduceTopicV3 {
1275                name: "orders".to_owned(),
1276                partitions: vec![ProducePartitionV3 {
1277                    partition_index: 0,
1278                    compression: RecordBatchCompression::None,
1279                    identity: RecordBatchIdentity {
1280                        producer_id: 42,
1281                        producer_epoch: 3,
1282                        base_sequence: 7,
1283                    },
1284                    records: vec![RecordBatchMessage::new(
1285                        Some(b"order-1".to_vec()),
1286                        Some(b"created".to_vec()),
1287                        1_000,
1288                    )],
1289                }],
1290            }],
1291        };
1292
1293        let bytes = request.encode().unwrap();
1294
1295        assert_eq!(&bytes[0..4], &[0, 0, 0, 9]);
1296        assert!(bytes
1297            .windows(7)
1298            .any(|window| window == [10, b'o', b'r', b'd', b'e', b'r', b's']));
1299        assert!(bytes.len() > 80);
1300    }
1301
1302    #[test]
1303    fn encodes_produce_request_v11_with_flexible_record_batch_schema() {
1304        let request = ProduceRequestV11 {
1305            correlation_id: 5,
1306            client_id: Some("kafrust".to_owned()),
1307            transactional_id: None,
1308            acks: 1,
1309            timeout_ms: 30_000,
1310            topics: Vec::new(),
1311        };
1312
1313        let bytes = request.encode().unwrap();
1314
1315        assert_eq!(&bytes[0..4], &[0, 0, 0, 11]);
1316        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1317        assert!(bytes.len() > 20);
1318    }
1319
1320    #[test]
1321    fn encodes_produce_request_v12_with_flexible_record_batch_schema() {
1322        let request = ProduceRequestV12 {
1323            correlation_id: 5,
1324            client_id: Some("kafrust".to_owned()),
1325            transactional_id: None,
1326            acks: 1,
1327            timeout_ms: 30_000,
1328            topics: Vec::new(),
1329        };
1330
1331        let bytes = request.encode().unwrap();
1332
1333        assert_eq!(&bytes[0..4], &[0, 0, 0, 12]);
1334        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1335        assert!(bytes.len() > 20);
1336    }
1337
1338    #[test]
1339    fn encodes_produce_request_v13_with_topic_uuid() {
1340        let request = ProduceRequestV13 {
1341            correlation_id: 5,
1342            client_id: Some("kafrust".to_owned()),
1343            transactional_id: None,
1344            acks: 1,
1345            timeout_ms: 30_000,
1346            topics: vec![ProduceTopicV13 {
1347                topic_id: [7; 16],
1348                partitions: Vec::new(),
1349            }],
1350        };
1351
1352        let bytes = request.encode().unwrap();
1353
1354        assert_eq!(&bytes[0..4], &[0, 0, 0, 13]);
1355        assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1356        assert!(bytes.windows(16).any(|window| window == [7; 16]));
1357        assert!(bytes.len() > 35);
1358    }
1359
1360    #[test]
1361    fn zstd_record_batch_encoding_roundtrips_through_fetch_decoder() {
1362        let record_set =
1363            encode_record_batch_set_with_compression(
1364                &[RecordBatchMessage::new(
1365                    Some(b"order-1".to_vec()),
1366                    Some(b"created".to_vec()),
1367                    1_000,
1368                )
1369                .header("source", Some(b"checkout".to_vec()))],
1370                RecordBatchCompression::Zstd,
1371            )
1372            .unwrap();
1373
1374        let mut bytes = Encoder::new();
1375        bytes.write_i32(0);
1376        bytes.write_i32(1);
1377        bytes.write_string("orders").unwrap();
1378        bytes.write_i32(1);
1379        bytes.write_i32(0);
1380        bytes.write_i16(0);
1381        bytes.write_i64(43);
1382        bytes.write_bytes(&record_set).unwrap();
1383        let bytes = bytes.into_bytes();
1384
1385        let mut decoder = Decoder::new(&bytes);
1386        let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1387        let record = &response.responses[0].partitions[0].records[0];
1388
1389        assert_eq!(record.offset, 0);
1390        assert_eq!(record.timestamp_ms, 1_000);
1391        assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1392        assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1393        assert!(decoder.is_empty());
1394    }
1395
1396    #[test]
1397    fn fetch_decoder_applies_custom_decompression_limit_to_record_batch() {
1398        let record_set = encode_record_batch_set_with_compression(
1399            &[RecordBatchMessage::new(
1400                Some(b"order-1".to_vec()),
1401                Some(vec![b'x'; 1024]),
1402                1_000,
1403            )],
1404            RecordBatchCompression::Zstd,
1405        )
1406        .unwrap();
1407
1408        let mut bytes = Encoder::new();
1409        bytes.write_i32(0);
1410        bytes.write_i32(1);
1411        bytes.write_string("orders").unwrap();
1412        bytes.write_i32(1);
1413        bytes.write_i32(0);
1414        bytes.write_i16(0);
1415        bytes.write_i64(43);
1416        bytes.write_bytes(&record_set).unwrap();
1417        let bytes = bytes.into_bytes();
1418        let limits = DecodeLimits::new().with_max_decompressed_record_bytes(64);
1419        let mut decoder = Decoder::with_limits(&bytes, limits);
1420
1421        assert!(matches!(
1422            FetchResponseV2::decode_body(&mut decoder),
1423            Err(crate::Error::LimitExceeded {
1424                kind: "decompressed record batch bytes",
1425                max: 64,
1426                ..
1427            })
1428        ));
1429    }
1430
1431    #[test]
1432    fn reports_message_set_encoded_len() {
1433        let records = [MessageSetMessage::new(
1434            Some(b"order-1".to_vec()),
1435            Some(b"created".to_vec()),
1436            1_000,
1437        )];
1438
1439        assert_eq!(
1440            encoded_message_set_len(&records).unwrap(),
1441            encode_message_set(&records).unwrap().len()
1442        );
1443    }
1444
1445    #[test]
1446    fn reports_record_batch_set_encoded_len() {
1447        let records =
1448            [
1449                RecordBatchMessage::new(
1450                    Some(b"order-1".to_vec()),
1451                    Some(b"created".to_vec()),
1452                    1_000,
1453                )
1454                .header("source", Some(b"checkout".to_vec())),
1455            ];
1456
1457        assert_eq!(
1458            encoded_record_batch_set_len(&records).unwrap(),
1459            encode_record_batch_set(&records).unwrap().len()
1460        );
1461    }
1462
1463    #[test]
1464    fn decodes_produce_response_v2() {
1465        let bytes = [
1466            0, 0, 0, 1, // topic response count
1467            0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
1468            0, 0, 0, 1, // partition response count
1469            0, 0, 0, 0, // partition
1470            0, 0, // error code
1471            0, 0, 0, 0, 0, 0, 0, 42, // base offset
1472            0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, // log append time -1
1473            0, 0, 0, 0, // throttle time
1474        ];
1475        let mut decoder = Decoder::new(&bytes);
1476        let response = ProduceResponseV2::decode_body(&mut decoder).unwrap();
1477
1478        assert_eq!(response.throttle_time_ms, 0);
1479        assert_eq!(response.responses[0].name, "orders");
1480        assert_eq!(response.responses[0].partitions[0].partition_index, 0);
1481        assert_eq!(response.responses[0].partitions[0].error_code, 0);
1482        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1483        assert!(decoder.is_empty());
1484    }
1485
1486    #[test]
1487    fn decodes_produce_response_v7() {
1488        let bytes = [
1489            0, 0, 0, 1, // topic response count
1490            0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic
1491            0, 0, 0, 1, // partition response count
1492            0, 0, 0, 0, // partition
1493            0, 0, // error code
1494            0, 0, 0, 0, 0, 0, 0, 42, // base offset
1495            0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, // log append time -1
1496            0, 0, 0, 0, 0, 0, 0, 7, // log start offset
1497            0, 0, 0, 0, // throttle time
1498        ];
1499        let mut decoder = Decoder::new(&bytes);
1500        let response = ProduceResponseV7::decode_body(&mut decoder).unwrap();
1501
1502        assert_eq!(response.throttle_time_ms, 0);
1503        assert_eq!(response.responses[0].name, "orders");
1504        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1505        assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1506        assert!(decoder.is_empty());
1507    }
1508
1509    #[test]
1510    fn decodes_produce_response_v9_with_record_error_fields() {
1511        let mut bytes = Encoder::new();
1512        bytes
1513            .write_compact_array(Some(&[()]), |encoder, _| {
1514                encoder.write_compact_string("orders")?;
1515                encoder.write_compact_array(Some(&[()]), |encoder, _| {
1516                    encoder.write_i32(0);
1517                    encoder.write_i16(0);
1518                    encoder.write_i64(42);
1519                    encoder.write_i64(-1);
1520                    encoder.write_i64(7);
1521                    encoder.write_compact_array(Some(&[()]), |encoder, _| {
1522                        encoder.write_i32(3);
1523                        encoder.write_compact_nullable_string(Some("bad record"))?;
1524                        encoder.write_empty_tagged_fields();
1525                        Ok(())
1526                    })?;
1527                    encoder.write_compact_nullable_string(Some("batch rejected"))?;
1528                    encoder.write_empty_tagged_fields();
1529                    Ok(())
1530                })?;
1531                encoder.write_empty_tagged_fields();
1532                Ok(())
1533            })
1534            .unwrap();
1535        bytes.write_i32(0);
1536        bytes.write_empty_tagged_fields();
1537
1538        let bytes = bytes.into_bytes();
1539        let mut decoder = Decoder::new(&bytes);
1540        let response = ProduceResponseV9::decode_body(&mut decoder).unwrap();
1541
1542        assert_eq!(response.responses[0].name, "orders");
1543        let partition = &response.responses[0].partitions[0];
1544        assert_eq!(partition.base_offset, 42);
1545        assert_eq!(partition.log_start_offset, 7);
1546        assert_eq!(partition.record_errors[0].batch_index, 3);
1547        assert_eq!(
1548            partition.record_errors[0].error_message.as_deref(),
1549            Some("bad record")
1550        );
1551        assert_eq!(partition.error_message.as_deref(), Some("batch rejected"));
1552        assert!(decoder.is_empty());
1553    }
1554
1555    #[test]
1556    fn decodes_produce_response_v13_with_topic_uuid() {
1557        let mut bytes = Encoder::new();
1558        bytes
1559            .write_compact_array(Some(&[()]), |encoder, _| {
1560                encoder.write_uuid(&[7; 16]);
1561                encoder.write_compact_array(Some(&[()]), |encoder, _| {
1562                    encoder.write_i32(0);
1563                    encoder.write_i16(0);
1564                    encoder.write_i64(42);
1565                    encoder.write_i64(-1);
1566                    encoder.write_i64(7);
1567                    encoder.write_unsigned_varint(1);
1568                    encoder.write_compact_nullable_string(None)?;
1569                    let mut current_leader = Encoder::new();
1570                    current_leader.write_i32(4);
1571                    current_leader.write_i32(12);
1572                    current_leader.write_empty_tagged_fields();
1573                    let current_leader = current_leader.into_bytes();
1574                    encoder.write_unsigned_varint(1);
1575                    encoder.write_unsigned_varint(0);
1576                    encoder.write_unsigned_varint(u32::try_from(current_leader.len()).unwrap());
1577                    encoder.write_raw(&current_leader);
1578                    Ok(())
1579                })?;
1580                encoder.write_empty_tagged_fields();
1581                Ok(())
1582            })
1583            .unwrap();
1584        bytes.write_i32(0);
1585        bytes.write_empty_tagged_fields();
1586
1587        let bytes = bytes.into_bytes();
1588        let mut decoder = Decoder::new(&bytes);
1589        let response = ProduceResponseV13::decode_body(&mut decoder).unwrap();
1590
1591        assert_eq!(response.responses[0].topic_id, [7; 16]);
1592        assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1593        assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1594        assert_eq!(
1595            response.responses[0].partitions[0].current_leader,
1596            Some(ProduceLeaderIdAndEpochV10 {
1597                leader_id: 4,
1598                leader_epoch: 12,
1599            })
1600        );
1601        assert!(response.node_endpoints.is_empty());
1602        assert!(decoder.is_empty());
1603    }
1604
1605    #[test]
1606    fn decodes_produce_response_node_endpoints_from_tag_zero() {
1607        let mut bytes = Encoder::new();
1608        bytes
1609            .write_compact_array::<()>(Some(&[]), |_, _| Ok(()))
1610            .unwrap();
1611        bytes.write_i32(0);
1612
1613        let mut endpoint = Encoder::new();
1614        endpoint
1615            .write_compact_array(Some(&[()]), |encoder, _| {
1616                encoder.write_i32(4);
1617                encoder.write_compact_string("broker-4")?;
1618                encoder.write_i32(9_092);
1619                encoder.write_compact_nullable_string(Some("rack-a"))?;
1620                encoder.write_empty_tagged_fields();
1621                Ok(())
1622            })
1623            .unwrap();
1624        let endpoint = endpoint.into_bytes();
1625        bytes.write_unsigned_varint(1);
1626        bytes.write_unsigned_varint(0);
1627        bytes.write_unsigned_varint(u32::try_from(endpoint.len()).unwrap());
1628        bytes.write_raw(&endpoint);
1629
1630        let bytes = bytes.into_bytes();
1631        let mut decoder = Decoder::new(&bytes);
1632        let response = ProduceResponseV13::decode_body(&mut decoder).unwrap();
1633
1634        assert_eq!(response.node_endpoints.len(), 1);
1635        assert_eq!(response.node_endpoints[0].node_id, 4);
1636        assert_eq!(response.node_endpoints[0].host, "broker-4");
1637        assert_eq!(response.node_endpoints[0].port, 9_092);
1638        assert_eq!(response.node_endpoints[0].rack.as_deref(), Some("rack-a"));
1639        assert!(decoder.is_empty());
1640    }
1641}