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