1use crate::codec::{Decoder, Encoder};
2use crate::error::{Error, Result};
3use crate::header::RequestHeader;
4use crate::record_batch::{compress_record_batch_records, RecordBatchCompression};
5
6pub const API_KEY: i16 = 0;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct ProduceRequestV2 {
10 pub correlation_id: i32,
11 pub client_id: Option<String>,
12 pub acks: i16,
13 pub timeout_ms: i32,
14 pub topics: Vec<ProduceTopicV2>,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct ProduceRequestV3 {
19 pub correlation_id: i32,
20 pub client_id: Option<String>,
21 pub transactional_id: Option<String>,
22 pub acks: i16,
23 pub timeout_ms: i32,
24 pub topics: Vec<ProduceTopicV3>,
25}
26
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct ProduceRequestV7 {
29 pub correlation_id: i32,
30 pub client_id: Option<String>,
31 pub transactional_id: Option<String>,
32 pub acks: i16,
33 pub timeout_ms: i32,
34 pub topics: Vec<ProduceTopicV3>,
35}
36
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct ProduceRequestV9 {
39 pub correlation_id: i32,
40 pub client_id: Option<String>,
41 pub transactional_id: Option<String>,
42 pub acks: i16,
43 pub timeout_ms: i32,
44 pub topics: Vec<ProduceTopicV3>,
45}
46
47#[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#[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#[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#[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
372pub fn encoded_message_set_len(records: &[MessageSetMessage]) -> Result<usize> {
374 Ok(encode_message_set(records)?.len())
375}
376
377pub fn encoded_record_batch_set_len(records: &[RecordBatchMessage]) -> Result<usize> {
379 Ok(encode_record_batch_set(records)?.len())
380}
381
382pub fn encoded_record_batch_set_len_with_compression(
384 records: &[RecordBatchMessage],
385 compression: RecordBatchCompression,
386) -> Result<usize> {
387 encoded_record_batch_set_len_with_compression_and_identity(
388 records,
389 compression,
390 RecordBatchIdentity::NON_IDEMPOTENT,
391 )
392}
393
394pub fn encoded_record_batch_set_len_with_compression_and_identity(
395 records: &[RecordBatchMessage],
396 compression: RecordBatchCompression,
397 identity: RecordBatchIdentity,
398) -> Result<usize> {
399 Ok(
400 encode_record_batch_set_with_compression_and_identity(records, compression, identity)?
401 .len(),
402 )
403}
404
405#[derive(Debug, Clone, PartialEq, Eq)]
406pub struct MessageSetMessage {
407 pub key: Option<Vec<u8>>,
408 pub value: Option<Vec<u8>>,
409 pub timestamp_ms: i64,
410}
411
412impl MessageSetMessage {
413 pub fn new(key: Option<Vec<u8>>, value: Option<Vec<u8>>, timestamp_ms: i64) -> Self {
414 Self {
415 key,
416 value,
417 timestamp_ms,
418 }
419 }
420}
421
422#[derive(Debug, Clone, PartialEq, Eq)]
423pub struct RecordBatchHeader {
424 pub key: String,
425 pub value: Option<Vec<u8>>,
426}
427
428impl RecordBatchHeader {
429 pub fn new(key: impl Into<String>, value: Option<Vec<u8>>) -> Self {
430 Self {
431 key: key.into(),
432 value,
433 }
434 }
435}
436
437#[derive(Debug, Clone, PartialEq, Eq)]
438pub struct RecordBatchMessage {
439 pub key: Option<Vec<u8>>,
440 pub value: Option<Vec<u8>>,
441 pub timestamp_ms: i64,
442 pub headers: Vec<RecordBatchHeader>,
443}
444
445impl RecordBatchMessage {
446 pub fn new(key: Option<Vec<u8>>, value: Option<Vec<u8>>, timestamp_ms: i64) -> Self {
447 Self {
448 key,
449 value,
450 timestamp_ms,
451 headers: Vec::new(),
452 }
453 }
454
455 pub fn header(mut self, key: impl Into<String>, value: Option<Vec<u8>>) -> Self {
456 self.headers.push(RecordBatchHeader::new(key, value));
457 self
458 }
459}
460
461#[derive(Debug, Clone, PartialEq, Eq)]
462pub struct ProduceResponseV2 {
463 pub responses: Vec<ProduceTopicResponseV2>,
464 pub throttle_time_ms: i32,
465}
466
467impl ProduceResponseV2 {
468 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
469 Ok(Self {
470 responses: decoder
471 .read_array("produce responses", ProduceTopicResponseV2::decode)?
472 .unwrap_or_default(),
473 throttle_time_ms: decoder.read_i32()?,
474 })
475 }
476}
477
478#[derive(Debug, Clone, PartialEq, Eq)]
479pub struct ProduceTopicResponseV2 {
480 pub name: String,
481 pub partitions: Vec<ProducePartitionResponseV2>,
482}
483
484impl ProduceTopicResponseV2 {
485 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
486 Ok(Self {
487 name: decoder.read_string()?,
488 partitions: decoder
489 .read_array(
490 "produce partition responses",
491 ProducePartitionResponseV2::decode,
492 )?
493 .unwrap_or_default(),
494 })
495 }
496}
497
498#[derive(Debug, Clone, PartialEq, Eq)]
499pub struct ProducePartitionResponseV2 {
500 pub partition_index: i32,
501 pub error_code: i16,
502 pub base_offset: i64,
503 pub log_append_time_ms: i64,
504}
505
506impl ProducePartitionResponseV2 {
507 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
508 Ok(Self {
509 partition_index: decoder.read_i32()?,
510 error_code: decoder.read_i16()?,
511 base_offset: decoder.read_i64()?,
512 log_append_time_ms: decoder.read_i64()?,
513 })
514 }
515}
516
517#[derive(Debug, Clone, PartialEq, Eq)]
518pub struct ProduceResponseV7 {
519 pub responses: Vec<ProduceTopicResponseV7>,
520 pub throttle_time_ms: i32,
521}
522
523#[derive(Debug, Clone, PartialEq, Eq)]
524pub struct ProduceResponseV9 {
525 pub responses: Vec<ProduceTopicResponseV9>,
526 pub throttle_time_ms: i32,
527}
528
529pub type ProduceResponseV11 = ProduceResponseV9;
531
532pub type ProduceResponseV12 = ProduceResponseV9;
534
535#[derive(Debug, Clone, PartialEq, Eq)]
537pub struct ProduceResponseV13 {
538 pub responses: Vec<ProduceTopicResponseV13>,
539 pub throttle_time_ms: i32,
540}
541
542impl ProduceResponseV13 {
543 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
544 let responses = decoder
545 .read_compact_array("produce responses", ProduceTopicResponseV13::decode)?
546 .unwrap_or_default();
547 let throttle_time_ms = decoder.read_i32()?;
548 decoder.read_tagged_fields()?;
549 Ok(Self {
550 responses,
551 throttle_time_ms,
552 })
553 }
554}
555
556#[derive(Debug, Clone, PartialEq, Eq)]
557pub struct ProduceTopicResponseV13 {
558 pub topic_id: [u8; 16],
559 pub partitions: Vec<ProducePartitionResponseV9>,
560}
561
562impl ProduceTopicResponseV13 {
563 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
564 let topic_id = decoder.read_uuid()?;
565 let partitions = decoder
566 .read_compact_array(
567 "produce partition responses",
568 ProducePartitionResponseV9::decode,
569 )?
570 .unwrap_or_default();
571 decoder.read_tagged_fields()?;
572 Ok(Self {
573 topic_id,
574 partitions,
575 })
576 }
577}
578
579impl ProduceResponseV9 {
580 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
581 let responses = decoder
582 .read_compact_array("produce responses", ProduceTopicResponseV9::decode)?
583 .unwrap_or_default();
584 let throttle_time_ms = decoder.read_i32()?;
585 decoder.read_tagged_fields()?;
586 Ok(Self {
587 responses,
588 throttle_time_ms,
589 })
590 }
591}
592
593#[derive(Debug, Clone, PartialEq, Eq)]
594pub struct ProduceTopicResponseV9 {
595 pub name: String,
596 pub partitions: Vec<ProducePartitionResponseV9>,
597}
598
599impl ProduceTopicResponseV9 {
600 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
601 let name = decoder.read_compact_string()?;
602 let partitions = decoder
603 .read_compact_array(
604 "produce partition responses",
605 ProducePartitionResponseV9::decode,
606 )?
607 .unwrap_or_default();
608 decoder.read_tagged_fields()?;
609 Ok(Self { name, partitions })
610 }
611}
612
613#[derive(Debug, Clone, PartialEq, Eq)]
614pub struct ProducePartitionResponseV9 {
615 pub partition_index: i32,
616 pub error_code: i16,
617 pub base_offset: i64,
618 pub log_append_time_ms: i64,
619 pub log_start_offset: i64,
620 pub record_errors: Vec<ProduceRecordErrorV9>,
621 pub error_message: Option<String>,
622}
623
624impl ProducePartitionResponseV9 {
625 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
626 let partition_index = decoder.read_i32()?;
627 let error_code = decoder.read_i16()?;
628 let base_offset = decoder.read_i64()?;
629 let log_append_time_ms = decoder.read_i64()?;
630 let log_start_offset = decoder.read_i64()?;
631 let record_errors = decoder
632 .read_compact_array("produce record errors", ProduceRecordErrorV9::decode)?
633 .unwrap_or_default();
634 let error_message = decoder.read_compact_nullable_string()?;
635 decoder.read_tagged_fields()?;
636 Ok(Self {
637 partition_index,
638 error_code,
639 base_offset,
640 log_append_time_ms,
641 log_start_offset,
642 record_errors,
643 error_message,
644 })
645 }
646}
647
648#[derive(Debug, Clone, PartialEq, Eq)]
649pub struct ProduceRecordErrorV9 {
650 pub batch_index: i32,
651 pub error_message: Option<String>,
652}
653
654impl ProduceRecordErrorV9 {
655 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
656 let batch_index = decoder.read_i32()?;
657 let error_message = decoder.read_compact_nullable_string()?;
658 decoder.read_tagged_fields()?;
659 Ok(Self {
660 batch_index,
661 error_message,
662 })
663 }
664}
665
666impl ProduceResponseV7 {
667 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
668 Ok(Self {
669 responses: decoder
670 .read_array("produce responses", ProduceTopicResponseV7::decode)?
671 .unwrap_or_default(),
672 throttle_time_ms: decoder.read_i32()?,
673 })
674 }
675}
676
677#[derive(Debug, Clone, PartialEq, Eq)]
678pub struct ProduceTopicResponseV7 {
679 pub name: String,
680 pub partitions: Vec<ProducePartitionResponseV7>,
681}
682
683impl ProduceTopicResponseV7 {
684 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
685 Ok(Self {
686 name: decoder.read_string()?,
687 partitions: decoder
688 .read_array(
689 "produce partition responses",
690 ProducePartitionResponseV7::decode,
691 )?
692 .unwrap_or_default(),
693 })
694 }
695}
696
697#[derive(Debug, Clone, PartialEq, Eq)]
698pub struct ProducePartitionResponseV7 {
699 pub partition_index: i32,
700 pub error_code: i16,
701 pub base_offset: i64,
702 pub log_append_time_ms: i64,
703 pub log_start_offset: i64,
704}
705
706impl ProducePartitionResponseV7 {
707 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
708 Ok(Self {
709 partition_index: decoder.read_i32()?,
710 error_code: decoder.read_i16()?,
711 base_offset: decoder.read_i64()?,
712 log_append_time_ms: decoder.read_i64()?,
713 log_start_offset: decoder.read_i64()?,
714 })
715 }
716}
717
718fn encode_message_set(records: &[MessageSetMessage]) -> Result<Vec<u8>> {
719 let mut set = Encoder::new();
720 for record in records {
721 let message = encode_message(record)?;
722 set.write_i64(0);
723 set.write_i32(i32::try_from(message.len()).map_err(|_| Error::LengthOverflow("message"))?);
724 set.write_raw(&message);
725 }
726 Ok(set.into_bytes())
727}
728
729fn encode_message(record: &MessageSetMessage) -> Result<Vec<u8>> {
730 let mut body = Encoder::new();
731 body.write_i8(1);
732 body.write_i8(0);
733 body.write_i64(record.timestamp_ms);
734 body.write_nullable_bytes(record.key.as_deref())?;
735 body.write_nullable_bytes(record.value.as_deref())?;
736 let body = body.into_bytes();
737
738 let mut message = Encoder::new();
739 message.write_i32(crc32_ieee(&body) as i32);
740 message.write_raw(&body);
741 Ok(message.into_bytes())
742}
743
744fn encode_record_batch_set(records: &[RecordBatchMessage]) -> Result<Vec<u8>> {
745 encode_record_batch_set_with_compression(records, RecordBatchCompression::None)
746}
747
748fn encode_record_batch_set_with_compression(
749 records: &[RecordBatchMessage],
750 compression: RecordBatchCompression,
751) -> Result<Vec<u8>> {
752 encode_record_batch_set_with_compression_and_identity(
753 records,
754 compression,
755 RecordBatchIdentity::NON_IDEMPOTENT,
756 )
757}
758
759fn encode_record_batch_set_with_compression_and_identity(
760 records: &[RecordBatchMessage],
761 compression: RecordBatchCompression,
762 identity: RecordBatchIdentity,
763) -> Result<Vec<u8>> {
764 encode_record_batch_set_with_compression_identity_and_transaction(
765 records,
766 compression,
767 identity,
768 false,
769 )
770}
771
772fn encode_record_batch_set_with_compression_identity_and_transaction(
773 records: &[RecordBatchMessage],
774 compression: RecordBatchCompression,
775 identity: RecordBatchIdentity,
776 transactional: bool,
777) -> Result<Vec<u8>> {
778 let base_timestamp = records
779 .first()
780 .map(|record| record.timestamp_ms)
781 .unwrap_or_default();
782 let max_timestamp = records
783 .iter()
784 .map(|record| record.timestamp_ms)
785 .max()
786 .unwrap_or(base_timestamp);
787 let last_offset_delta = records
788 .len()
789 .checked_sub(1)
790 .map(|delta| i32::try_from(delta).map_err(|_| Error::LengthOverflow("record batch")))
791 .transpose()?
792 .unwrap_or_default();
793
794 let record_count =
795 i32::try_from(records.len()).map_err(|_| Error::LengthOverflow("record batch records"))?;
796 let mut record_bytes = Encoder::new();
797 for (offset_delta, record) in records.iter().enumerate() {
798 let encoded = encode_record(record, base_timestamp, offset_delta)?;
799 record_bytes.write_varint(
800 i32::try_from(encoded.len()).map_err(|_| Error::LengthOverflow("record"))?,
801 );
802 record_bytes.write_raw(&encoded);
803 }
804 let record_bytes = compress_record_batch_records(compression, &record_bytes.into_bytes())?;
805
806 let mut crc_payload = Encoder::new();
807 let attributes = compression.attributes() | if transactional { 0x10 } else { 0 };
808 crc_payload.write_i16(attributes);
809 crc_payload.write_i32(last_offset_delta);
810 crc_payload.write_i64(base_timestamp);
811 crc_payload.write_i64(max_timestamp);
812 crc_payload.write_i64(identity.producer_id);
813 crc_payload.write_i16(identity.producer_epoch);
814 crc_payload.write_i32(identity.base_sequence);
815 crc_payload.write_i32(record_count);
816 crc_payload.write_raw(&record_bytes);
817 let crc_payload = crc_payload.into_bytes();
818
819 let mut batch = Encoder::new();
820 batch.write_i32(0);
821 batch.write_i8(2);
822 batch.write_i32(crc32c(&crc_payload) as i32);
823 batch.write_raw(&crc_payload);
824 let batch = batch.into_bytes();
825
826 let mut set = Encoder::new();
827 set.write_i64(0);
828 set.write_i32(i32::try_from(batch.len()).map_err(|_| Error::LengthOverflow("record batch"))?);
829 set.write_raw(&batch);
830 Ok(set.into_bytes())
831}
832
833fn encode_record(
834 record: &RecordBatchMessage,
835 base_timestamp: i64,
836 offset_delta: usize,
837) -> Result<Vec<u8>> {
838 let mut encoder = Encoder::new();
839 encoder.write_i8(0);
840 encoder.write_varlong(record.timestamp_ms.saturating_sub(base_timestamp));
841 encoder.write_varint(
842 i32::try_from(offset_delta).map_err(|_| Error::LengthOverflow("record offset delta"))?,
843 );
844 encoder.write_varint_nullable_bytes(record.key.as_deref())?;
845 encoder.write_varint_nullable_bytes(record.value.as_deref())?;
846 encoder.write_varint(
847 i32::try_from(record.headers.len()).map_err(|_| Error::LengthOverflow("record headers"))?,
848 );
849 for header in &record.headers {
850 encoder.write_varint_bytes(header.key.as_bytes())?;
851 encoder.write_varint_nullable_bytes(header.value.as_deref())?;
852 }
853 Ok(encoder.into_bytes())
854}
855
856fn crc32_ieee(bytes: &[u8]) -> u32 {
857 crc32_with_table(bytes, &CRC32_IEEE_TABLE)
858}
859
860fn crc32c(bytes: &[u8]) -> u32 {
861 crc32_with_table(bytes, &CRC32C_TABLE)
862}
863
864const CRC32_IEEE_TABLE: [u32; 256] = crc32_table(0xedb8_8320);
865const CRC32C_TABLE: [u32; 256] = crc32_table(0x82f6_3b78);
866
867const fn crc32_table(polynomial: u32) -> [u32; 256] {
868 let mut table = [0; 256];
869 let mut index = 0;
870 while index < table.len() {
871 let mut value = index as u32;
872 let mut bit = 0;
873 while bit < 8 {
874 let mask = 0u32.wrapping_sub(value & 1);
875 value = (value >> 1) ^ (polynomial & mask);
876 bit += 1;
877 }
878 table[index] = value;
879 index += 1;
880 }
881 table
882}
883
884fn crc32_with_table(bytes: &[u8], table: &[u32; 256]) -> u32 {
885 let mut crc = 0xffff_ffffu32;
886 for byte in bytes {
887 let index = usize::from((crc as u8) ^ byte);
888 crc = (crc >> 8) ^ table[index];
889 }
890 !crc
891}
892
893#[cfg(test)]
894#[allow(clippy::unwrap_used)]
895mod tests {
896 use super::{
897 crc32_ieee, crc32c, encode_message_set, encode_record_batch_set,
898 encode_record_batch_set_with_compression,
899 encode_record_batch_set_with_compression_and_identity,
900 encode_record_batch_set_with_compression_identity_and_transaction, encoded_message_set_len,
901 encoded_record_batch_set_len, MessageSetMessage, ProducePartitionV2, ProducePartitionV3,
902 ProduceRequestV11, ProduceRequestV12, ProduceRequestV13, ProduceRequestV2,
903 ProduceRequestV3, ProduceRequestV7, ProduceRequestV9, ProduceResponseV13,
904 ProduceResponseV2, ProduceResponseV7, ProduceResponseV9, ProduceTopicV13, ProduceTopicV2,
905 ProduceTopicV3, RecordBatchIdentity, RecordBatchMessage,
906 };
907 use crate::codec::{DecodeLimits, Decoder};
908 use crate::record_batch::RecordBatchCompression;
909 use crate::{api::fetch::FetchResponseV2, codec::Encoder};
910
911 #[test]
912 fn crc_implementations_match_standard_check_vectors() {
913 assert_eq!(crc32_ieee(b"123456789"), 0xcbf4_3926);
914 assert_eq!(crc32c(b"123456789"), 0xe306_9283);
915 }
916
917 #[test]
918 fn encodes_produce_request_v2() {
919 let request = ProduceRequestV2 {
920 correlation_id: 5,
921 client_id: Some("kafrust".to_owned()),
922 acks: 1,
923 timeout_ms: 30_000,
924 topics: vec![ProduceTopicV2 {
925 name: "orders".to_owned(),
926 partitions: vec![ProducePartitionV2 {
927 partition_index: 0,
928 records: vec![MessageSetMessage::new(
929 Some(b"order-1".to_vec()),
930 Some(b"created".to_vec()),
931 0,
932 )],
933 }],
934 }],
935 };
936
937 let bytes = request.encode().unwrap();
938 assert_eq!(
939 &bytes[0..17],
940 &[0, 0, 0, 2, 0, 0, 0, 5, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',]
941 );
942 assert!(bytes.len() > 60);
943 }
944
945 #[test]
946 fn encodes_idempotent_record_batch_identity() {
947 let set = encode_record_batch_set_with_compression_and_identity(
948 &[RecordBatchMessage::new(
949 Some(b"order-1".to_vec()),
950 Some(b"created".to_vec()),
951 1_000,
952 )],
953 RecordBatchCompression::None,
954 RecordBatchIdentity {
955 producer_id: 42,
956 producer_epoch: 3,
957 base_sequence: 7,
958 },
959 )
960 .unwrap();
961
962 assert_eq!(&set[43..51], &42_i64.to_be_bytes());
963 assert_eq!(&set[51..53], &3_i16.to_be_bytes());
964 assert_eq!(&set[53..57], &7_i32.to_be_bytes());
965 }
966
967 #[test]
968 fn encodes_transactional_record_batch_attribute() {
969 let set = encode_record_batch_set_with_compression_identity_and_transaction(
970 &[RecordBatchMessage::new(
971 None,
972 Some(b"created".to_vec()),
973 1_000,
974 )],
975 RecordBatchCompression::None,
976 RecordBatchIdentity {
977 producer_id: 42,
978 producer_epoch: 3,
979 base_sequence: 0,
980 },
981 true,
982 )
983 .unwrap();
984
985 assert_eq!(&set[21..23], &0x10_i16.to_be_bytes());
986 }
987
988 #[test]
989 fn encodes_produce_request_v3_with_record_batch() {
990 let request = ProduceRequestV3 {
991 correlation_id: 5,
992 client_id: Some("kafrust".to_owned()),
993 transactional_id: None,
994 acks: 1,
995 timeout_ms: 30_000,
996 topics: vec![ProduceTopicV3 {
997 name: "orders".to_owned(),
998 partitions: vec![ProducePartitionV3 {
999 partition_index: 0,
1000 compression: RecordBatchCompression::None,
1001 identity: RecordBatchIdentity::NON_IDEMPOTENT,
1002 records: vec![RecordBatchMessage::new(
1003 Some(b"order-1".to_vec()),
1004 Some(b"created".to_vec()),
1005 1_000,
1006 )
1007 .header("source", Some(b"checkout".to_vec()))],
1008 }],
1009 }],
1010 };
1011
1012 let bytes = request.encode().unwrap();
1013
1014 assert_eq!(&bytes[0..4], &[0, 0, 0, 3]);
1015 assert!(bytes.len() > 80);
1016 }
1017
1018 #[test]
1019 fn record_batch_encoding_roundtrips_through_fetch_decoder() {
1020 let record_set = encode_record_batch_set(&[RecordBatchMessage::new(
1021 Some(b"order-1".to_vec()),
1022 Some(b"created".to_vec()),
1023 1_000,
1024 )
1025 .header("source", Some(b"checkout".to_vec()))])
1026 .unwrap();
1027
1028 let mut bytes = Encoder::new();
1029 bytes.write_i32(0);
1030 bytes.write_i32(1);
1031 bytes.write_string("orders").unwrap();
1032 bytes.write_i32(1);
1033 bytes.write_i32(0);
1034 bytes.write_i16(0);
1035 bytes.write_i64(43);
1036 bytes.write_bytes(&record_set).unwrap();
1037 let bytes = bytes.into_bytes();
1038
1039 let mut decoder = Decoder::new(&bytes);
1040 let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1041 let record = &response.responses[0].partitions[0].records[0];
1042
1043 assert_eq!(record.offset, 0);
1044 assert_eq!(record.timestamp_ms, 1_000);
1045 assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1046 assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1047 assert!(decoder.is_empty());
1048 }
1049
1050 #[test]
1051 fn gzip_record_batch_encoding_roundtrips_through_fetch_decoder() {
1052 let record_set =
1053 encode_record_batch_set_with_compression(
1054 &[RecordBatchMessage::new(
1055 Some(b"order-1".to_vec()),
1056 Some(b"created".to_vec()),
1057 1_000,
1058 )
1059 .header("source", Some(b"checkout".to_vec()))],
1060 RecordBatchCompression::Gzip,
1061 )
1062 .unwrap();
1063
1064 let mut bytes = Encoder::new();
1065 bytes.write_i32(0);
1066 bytes.write_i32(1);
1067 bytes.write_string("orders").unwrap();
1068 bytes.write_i32(1);
1069 bytes.write_i32(0);
1070 bytes.write_i16(0);
1071 bytes.write_i64(43);
1072 bytes.write_bytes(&record_set).unwrap();
1073 let bytes = bytes.into_bytes();
1074
1075 let mut decoder = Decoder::new(&bytes);
1076 let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1077 let record = &response.responses[0].partitions[0].records[0];
1078
1079 assert_eq!(record.offset, 0);
1080 assert_eq!(record.timestamp_ms, 1_000);
1081 assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1082 assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1083 assert!(decoder.is_empty());
1084 }
1085
1086 #[test]
1087 fn snappy_record_batch_encoding_roundtrips_through_fetch_decoder() {
1088 let record_set =
1089 encode_record_batch_set_with_compression(
1090 &[RecordBatchMessage::new(
1091 Some(b"order-1".to_vec()),
1092 Some(b"created".to_vec()),
1093 1_000,
1094 )
1095 .header("source", Some(b"checkout".to_vec()))],
1096 RecordBatchCompression::Snappy,
1097 )
1098 .unwrap();
1099
1100 let mut bytes = Encoder::new();
1101 bytes.write_i32(0);
1102 bytes.write_i32(1);
1103 bytes.write_string("orders").unwrap();
1104 bytes.write_i32(1);
1105 bytes.write_i32(0);
1106 bytes.write_i16(0);
1107 bytes.write_i64(43);
1108 bytes.write_bytes(&record_set).unwrap();
1109 let bytes = bytes.into_bytes();
1110
1111 let mut decoder = Decoder::new(&bytes);
1112 let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1113 let record = &response.responses[0].partitions[0].records[0];
1114
1115 assert_eq!(record.offset, 0);
1116 assert_eq!(record.timestamp_ms, 1_000);
1117 assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1118 assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1119 assert!(decoder.is_empty());
1120 }
1121
1122 #[test]
1123 fn lz4_record_batch_encoding_roundtrips_through_fetch_decoder() {
1124 let record_set =
1125 encode_record_batch_set_with_compression(
1126 &[RecordBatchMessage::new(
1127 Some(b"order-1".to_vec()),
1128 Some(b"created".to_vec()),
1129 1_000,
1130 )
1131 .header("source", Some(b"checkout".to_vec()))],
1132 RecordBatchCompression::Lz4,
1133 )
1134 .unwrap();
1135
1136 let mut bytes = Encoder::new();
1137 bytes.write_i32(0);
1138 bytes.write_i32(1);
1139 bytes.write_string("orders").unwrap();
1140 bytes.write_i32(1);
1141 bytes.write_i32(0);
1142 bytes.write_i16(0);
1143 bytes.write_i64(43);
1144 bytes.write_bytes(&record_set).unwrap();
1145 let bytes = bytes.into_bytes();
1146
1147 let mut decoder = Decoder::new(&bytes);
1148 let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1149 let record = &response.responses[0].partitions[0].records[0];
1150
1151 assert_eq!(record.offset, 0);
1152 assert_eq!(record.timestamp_ms, 1_000);
1153 assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1154 assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1155 assert!(decoder.is_empty());
1156 }
1157
1158 #[test]
1159 fn encodes_produce_request_v7_with_record_batch() {
1160 let request = ProduceRequestV7 {
1161 correlation_id: 5,
1162 client_id: Some("kafrust".to_owned()),
1163 transactional_id: None,
1164 acks: 1,
1165 timeout_ms: 30_000,
1166 topics: vec![ProduceTopicV3 {
1167 name: "orders".to_owned(),
1168 partitions: vec![ProducePartitionV3 {
1169 partition_index: 0,
1170 compression: RecordBatchCompression::Zstd,
1171 identity: RecordBatchIdentity::NON_IDEMPOTENT,
1172 records: vec![RecordBatchMessage::new(
1173 Some(b"order-1".to_vec()),
1174 Some(b"created".to_vec()),
1175 1_000,
1176 )],
1177 }],
1178 }],
1179 };
1180
1181 let bytes = request.encode().unwrap();
1182
1183 assert_eq!(&bytes[0..4], &[0, 0, 0, 7]);
1184 assert!(bytes.len() > 70);
1185 }
1186
1187 #[test]
1188 fn encodes_produce_request_v9_with_flexible_record_batch() {
1189 let request = ProduceRequestV9 {
1190 correlation_id: 5,
1191 client_id: Some("kafrust".to_owned()),
1192 transactional_id: Some("orders-tx".to_owned()),
1193 acks: -1,
1194 timeout_ms: 30_000,
1195 topics: vec![ProduceTopicV3 {
1196 name: "orders".to_owned(),
1197 partitions: vec![ProducePartitionV3 {
1198 partition_index: 0,
1199 compression: RecordBatchCompression::None,
1200 identity: RecordBatchIdentity {
1201 producer_id: 42,
1202 producer_epoch: 3,
1203 base_sequence: 7,
1204 },
1205 records: vec![RecordBatchMessage::new(
1206 Some(b"order-1".to_vec()),
1207 Some(b"created".to_vec()),
1208 1_000,
1209 )],
1210 }],
1211 }],
1212 };
1213
1214 let bytes = request.encode().unwrap();
1215
1216 assert_eq!(&bytes[0..4], &[0, 0, 0, 9]);
1217 assert!(bytes
1218 .windows(7)
1219 .any(|window| window == [10, b'o', b'r', b'd', b'e', b'r', b's']));
1220 assert!(bytes.len() > 80);
1221 }
1222
1223 #[test]
1224 fn encodes_produce_request_v11_with_flexible_record_batch_schema() {
1225 let request = ProduceRequestV11 {
1226 correlation_id: 5,
1227 client_id: Some("kafrust".to_owned()),
1228 transactional_id: None,
1229 acks: 1,
1230 timeout_ms: 30_000,
1231 topics: Vec::new(),
1232 };
1233
1234 let bytes = request.encode().unwrap();
1235
1236 assert_eq!(&bytes[0..4], &[0, 0, 0, 11]);
1237 assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1238 assert!(bytes.len() > 20);
1239 }
1240
1241 #[test]
1242 fn encodes_produce_request_v12_with_flexible_record_batch_schema() {
1243 let request = ProduceRequestV12 {
1244 correlation_id: 5,
1245 client_id: Some("kafrust".to_owned()),
1246 transactional_id: None,
1247 acks: 1,
1248 timeout_ms: 30_000,
1249 topics: Vec::new(),
1250 };
1251
1252 let bytes = request.encode().unwrap();
1253
1254 assert_eq!(&bytes[0..4], &[0, 0, 0, 12]);
1255 assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1256 assert!(bytes.len() > 20);
1257 }
1258
1259 #[test]
1260 fn encodes_produce_request_v13_with_topic_uuid() {
1261 let request = ProduceRequestV13 {
1262 correlation_id: 5,
1263 client_id: Some("kafrust".to_owned()),
1264 transactional_id: None,
1265 acks: 1,
1266 timeout_ms: 30_000,
1267 topics: vec![ProduceTopicV13 {
1268 topic_id: [7; 16],
1269 partitions: Vec::new(),
1270 }],
1271 };
1272
1273 let bytes = request.encode().unwrap();
1274
1275 assert_eq!(&bytes[0..4], &[0, 0, 0, 13]);
1276 assert_eq!(&bytes[4..8], &[0, 0, 0, 5]);
1277 assert!(bytes.windows(16).any(|window| window == [7; 16]));
1278 assert!(bytes.len() > 35);
1279 }
1280
1281 #[test]
1282 fn zstd_record_batch_encoding_roundtrips_through_fetch_decoder() {
1283 let record_set =
1284 encode_record_batch_set_with_compression(
1285 &[RecordBatchMessage::new(
1286 Some(b"order-1".to_vec()),
1287 Some(b"created".to_vec()),
1288 1_000,
1289 )
1290 .header("source", Some(b"checkout".to_vec()))],
1291 RecordBatchCompression::Zstd,
1292 )
1293 .unwrap();
1294
1295 let mut bytes = Encoder::new();
1296 bytes.write_i32(0);
1297 bytes.write_i32(1);
1298 bytes.write_string("orders").unwrap();
1299 bytes.write_i32(1);
1300 bytes.write_i32(0);
1301 bytes.write_i16(0);
1302 bytes.write_i64(43);
1303 bytes.write_bytes(&record_set).unwrap();
1304 let bytes = bytes.into_bytes();
1305
1306 let mut decoder = Decoder::new(&bytes);
1307 let response = FetchResponseV2::decode_body(&mut decoder).unwrap();
1308 let record = &response.responses[0].partitions[0].records[0];
1309
1310 assert_eq!(record.offset, 0);
1311 assert_eq!(record.timestamp_ms, 1_000);
1312 assert_eq!(record.key.as_deref(), Some(&b"order-1"[..]));
1313 assert_eq!(record.value.as_deref(), Some(&b"created"[..]));
1314 assert!(decoder.is_empty());
1315 }
1316
1317 #[test]
1318 fn fetch_decoder_applies_custom_decompression_limit_to_record_batch() {
1319 let record_set = encode_record_batch_set_with_compression(
1320 &[RecordBatchMessage::new(
1321 Some(b"order-1".to_vec()),
1322 Some(vec![b'x'; 1024]),
1323 1_000,
1324 )],
1325 RecordBatchCompression::Zstd,
1326 )
1327 .unwrap();
1328
1329 let mut bytes = Encoder::new();
1330 bytes.write_i32(0);
1331 bytes.write_i32(1);
1332 bytes.write_string("orders").unwrap();
1333 bytes.write_i32(1);
1334 bytes.write_i32(0);
1335 bytes.write_i16(0);
1336 bytes.write_i64(43);
1337 bytes.write_bytes(&record_set).unwrap();
1338 let bytes = bytes.into_bytes();
1339 let limits = DecodeLimits::new().with_max_decompressed_record_bytes(64);
1340 let mut decoder = Decoder::with_limits(&bytes, limits);
1341
1342 assert!(matches!(
1343 FetchResponseV2::decode_body(&mut decoder),
1344 Err(crate::Error::LimitExceeded {
1345 kind: "decompressed record batch bytes",
1346 max: 64,
1347 ..
1348 })
1349 ));
1350 }
1351
1352 #[test]
1353 fn reports_message_set_encoded_len() {
1354 let records = [MessageSetMessage::new(
1355 Some(b"order-1".to_vec()),
1356 Some(b"created".to_vec()),
1357 1_000,
1358 )];
1359
1360 assert_eq!(
1361 encoded_message_set_len(&records).unwrap(),
1362 encode_message_set(&records).unwrap().len()
1363 );
1364 }
1365
1366 #[test]
1367 fn reports_record_batch_set_encoded_len() {
1368 let records =
1369 [
1370 RecordBatchMessage::new(
1371 Some(b"order-1".to_vec()),
1372 Some(b"created".to_vec()),
1373 1_000,
1374 )
1375 .header("source", Some(b"checkout".to_vec())),
1376 ];
1377
1378 assert_eq!(
1379 encoded_record_batch_set_len(&records).unwrap(),
1380 encode_record_batch_set(&records).unwrap().len()
1381 );
1382 }
1383
1384 #[test]
1385 fn decodes_produce_response_v2() {
1386 let bytes = [
1387 0, 0, 0, 1, 0, 6, b'o', b'r', b'd', b'e', b'r', b's', 0, 0, 0, 1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 42, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0, 0, 0, 0, ];
1396 let mut decoder = Decoder::new(&bytes);
1397 let response = ProduceResponseV2::decode_body(&mut decoder).unwrap();
1398
1399 assert_eq!(response.throttle_time_ms, 0);
1400 assert_eq!(response.responses[0].name, "orders");
1401 assert_eq!(response.responses[0].partitions[0].partition_index, 0);
1402 assert_eq!(response.responses[0].partitions[0].error_code, 0);
1403 assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1404 assert!(decoder.is_empty());
1405 }
1406
1407 #[test]
1408 fn decodes_produce_response_v7() {
1409 let bytes = [
1410 0, 0, 0, 1, 0, 6, b'o', b'r', b'd', b'e', b'r', b's', 0, 0, 0, 1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 42, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0, 0, 0, 0, 0, 0, 0, 7, 0, 0, 0, 0, ];
1420 let mut decoder = Decoder::new(&bytes);
1421 let response = ProduceResponseV7::decode_body(&mut decoder).unwrap();
1422
1423 assert_eq!(response.throttle_time_ms, 0);
1424 assert_eq!(response.responses[0].name, "orders");
1425 assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1426 assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1427 assert!(decoder.is_empty());
1428 }
1429
1430 #[test]
1431 fn decodes_produce_response_v9_with_record_error_fields() {
1432 let mut bytes = Encoder::new();
1433 bytes
1434 .write_compact_array(Some(&[()]), |encoder, _| {
1435 encoder.write_compact_string("orders")?;
1436 encoder.write_compact_array(Some(&[()]), |encoder, _| {
1437 encoder.write_i32(0);
1438 encoder.write_i16(0);
1439 encoder.write_i64(42);
1440 encoder.write_i64(-1);
1441 encoder.write_i64(7);
1442 encoder.write_compact_array(Some(&[()]), |encoder, _| {
1443 encoder.write_i32(3);
1444 encoder.write_compact_nullable_string(Some("bad record"))?;
1445 encoder.write_empty_tagged_fields();
1446 Ok(())
1447 })?;
1448 encoder.write_compact_nullable_string(Some("batch rejected"))?;
1449 encoder.write_empty_tagged_fields();
1450 Ok(())
1451 })?;
1452 encoder.write_empty_tagged_fields();
1453 Ok(())
1454 })
1455 .unwrap();
1456 bytes.write_i32(0);
1457 bytes.write_empty_tagged_fields();
1458
1459 let bytes = bytes.into_bytes();
1460 let mut decoder = Decoder::new(&bytes);
1461 let response = ProduceResponseV9::decode_body(&mut decoder).unwrap();
1462
1463 assert_eq!(response.responses[0].name, "orders");
1464 let partition = &response.responses[0].partitions[0];
1465 assert_eq!(partition.base_offset, 42);
1466 assert_eq!(partition.log_start_offset, 7);
1467 assert_eq!(partition.record_errors[0].batch_index, 3);
1468 assert_eq!(
1469 partition.record_errors[0].error_message.as_deref(),
1470 Some("bad record")
1471 );
1472 assert_eq!(partition.error_message.as_deref(), Some("batch rejected"));
1473 assert!(decoder.is_empty());
1474 }
1475
1476 #[test]
1477 fn decodes_produce_response_v13_with_topic_uuid() {
1478 let mut bytes = Encoder::new();
1479 bytes
1480 .write_compact_array(Some(&[()]), |encoder, _| {
1481 encoder.write_uuid(&[7; 16]);
1482 encoder.write_compact_array(Some(&[()]), |encoder, _| {
1483 encoder.write_i32(0);
1484 encoder.write_i16(0);
1485 encoder.write_i64(42);
1486 encoder.write_i64(-1);
1487 encoder.write_i64(7);
1488 encoder.write_unsigned_varint(1);
1489 encoder.write_compact_nullable_string(None)?;
1490 encoder.write_empty_tagged_fields();
1491 Ok(())
1492 })?;
1493 encoder.write_empty_tagged_fields();
1494 Ok(())
1495 })
1496 .unwrap();
1497 bytes.write_i32(0);
1498 bytes.write_empty_tagged_fields();
1499
1500 let bytes = bytes.into_bytes();
1501 let mut decoder = Decoder::new(&bytes);
1502 let response = ProduceResponseV13::decode_body(&mut decoder).unwrap();
1503
1504 assert_eq!(response.responses[0].topic_id, [7; 16]);
1505 assert_eq!(response.responses[0].partitions[0].base_offset, 42);
1506 assert_eq!(response.responses[0].partitions[0].log_start_offset, 7);
1507 assert!(decoder.is_empty());
1508 }
1509}