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#[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 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 pub node_endpoints: Vec<ProduceNodeEndpointV10>,
531}
532
533pub type ProduceResponseV11 = ProduceResponseV9;
535
536pub type ProduceResponseV12 = ProduceResponseV9;
538
539#[derive(Debug, Clone, PartialEq, Eq)]
541pub struct ProduceResponseV13 {
542 pub responses: Vec<ProduceTopicResponseV13>,
543 pub throttle_time_ms: i32,
544 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
577pub struct ProduceLeaderIdAndEpochV10 {
578 pub leader_id: i32,
579 pub leader_epoch: i32,
580}
581
582#[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 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, 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, ];
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, 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, ];
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(¤t_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}