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