1use std::sync::Arc;
2
3use arrow::array::ArrayData;
4use arrow::buffer::Buffer;
5use arrow_array::builder::ArrayBuilder;
6use arrow_array::builder::{
7 BinaryBuilder, BooleanBuilder, LargeBinaryBuilder, LargeStringBuilder, PrimitiveBuilder,
8 StringBuilder,
9};
10use arrow_array::types::{
11 Date32Type, Float32Type, Float64Type, Int32Type, Int64Type, UInt32Type, UInt64Type,
12};
13use arrow_array::{Array, ArrowPrimitiveType, BinaryArray, MapArray, RecordBatch, StructArray};
14use arrow_schema::{DataType, Field, TimeUnit};
15use chrono::Datelike;
16
17use prost_reflect::{EnumDescriptor, FieldDescriptor, Kind, MessageDescriptor};
18
19use crate::config::{ConfluentWirePolicy, EnumRepr, PtarsConfig};
20
21#[allow(deprecated)]
26fn decode_error(msg: &str) -> prost::DecodeError {
27 prost::DecodeError::new(msg.to_string())
28}
29
30fn decode_varint(buf: &[u8]) -> Result<(u64, usize), prost::DecodeError> {
31 let mut result: u64 = 0;
32 let mut shift = 0u32;
33 for (i, &byte) in buf.iter().enumerate() {
34 result |= ((byte & 0x7F) as u64) << shift;
35 if byte & 0x80 == 0 {
36 return Ok((result, i + 1));
37 }
38 shift += 7;
39 if shift >= 64 {
40 return Err(decode_error("varint too large"));
41 }
42 }
43 Err(decode_error("unexpected EOF in varint"))
44}
45
46fn decode_tag(buf: &[u8]) -> Result<(u32, u8, usize), prost::DecodeError> {
47 let (key, n) = decode_varint(buf)?;
48 let wire_type = (key & 0x07) as u8;
49 let field_number = (key >> 3) as u32;
50 if field_number == 0 {
51 return Err(decode_error("invalid field number 0"));
52 }
53 Ok((field_number, wire_type, n))
54}
55
56fn skip_field(wire_type: u8, buf: &[u8]) -> Result<usize, prost::DecodeError> {
57 match wire_type {
58 0 => {
59 let (_, n) = decode_varint(buf)?;
60 Ok(n)
61 }
62 1 => {
63 if buf.len() < 8 {
64 return Err(decode_error("unexpected EOF"));
65 }
66 Ok(8)
67 }
68 2 => {
69 let (len, n) = decode_varint(buf)?;
70 let total = n + len as usize;
71 if buf.len() < total {
72 return Err(decode_error("unexpected EOF"));
73 }
74 Ok(total)
75 }
76 5 => {
77 if buf.len() < 4 {
78 return Err(decode_error("unexpected EOF"));
79 }
80 Ok(4)
81 }
82 _ => Err(decode_error("unsupported wire type")),
83 }
84}
85
86fn read_length_delimited(buf: &[u8]) -> Result<(&[u8], usize), prost::DecodeError> {
88 let (len, n) = decode_varint(buf)?;
89 let len = len as usize;
90 let total = n + len;
91 if buf.len() < total {
92 return Err(decode_error("unexpected EOF"));
93 }
94 Ok((&buf[n..total], total))
95}
96
97fn strip_confluent_prefix(
99 buf: &[u8],
100 policy: ConfluentWirePolicy,
101) -> Result<&[u8], prost::DecodeError> {
102 match policy {
103 ConfluentWirePolicy::Raw => Ok(buf),
104 ConfluentWirePolicy::Standard => {
105 if buf.len() < 5 {
106 return Err(decode_error(
107 "message too short for Confluent wire format header",
108 ));
109 }
110 Ok(&buf[5..])
111 }
112 ConfluentWirePolicy::Protobuf => {
113 if buf.len() < 5 {
114 return Err(decode_error(
115 "message too short for Confluent wire format header",
116 ));
117 }
118 let remaining = &buf[5..];
119 let (count, mut offset) = decode_varint(remaining)?;
121 for _ in 0..count {
123 let (_, n) = decode_varint(&remaining[offset..])?;
124 offset += n;
125 }
126 Ok(&remaining[offset..])
127 }
128 }
129}
130
131#[inline]
132fn decode_zigzag32(v: u64) -> i32 {
133 let v = v as u32;
134 ((v >> 1) as i32) ^ (-((v & 1) as i32))
135}
136
137#[inline]
138fn decode_zigzag64(v: u64) -> i64 {
139 ((v >> 1) as i64) ^ (-((v & 1) as i64))
140}
141
142fn convert_seconds_nanos_to_unit(seconds: i64, nanos: i32, unit: TimeUnit, type_name: &str) -> i64 {
143 match unit {
144 TimeUnit::Second => seconds,
145 TimeUnit::Millisecond => seconds
146 .checked_mul(1_000)
147 .and_then(|s| s.checked_add(i64::from(nanos) / 1_000_000))
148 .unwrap_or_else(|| panic!("{type_name} overflow")),
149 TimeUnit::Microsecond => seconds
150 .checked_mul(1_000_000)
151 .and_then(|s| s.checked_add(i64::from(nanos) / 1_000))
152 .unwrap_or_else(|| panic!("{type_name} overflow")),
153 TimeUnit::Nanosecond => seconds
154 .checked_mul(1_000_000_000)
155 .and_then(|s| s.checked_add(i64::from(nanos)))
156 .unwrap_or_else(|| panic!("{type_name} overflow")),
157 }
158}
159
160static CE_OFFSET: i32 = 719163;
161
162fn enum_name(enum_descriptor: &EnumDescriptor, number: i32) -> String {
163 match enum_descriptor.get_value(number) {
164 Some(v) => v.name().to_string(),
165 None => number.to_string(),
166 }
167}
168
169enum StringBuilderInner {
174 Regular(StringBuilder),
175 Large(LargeStringBuilder),
176}
177
178impl StringBuilderInner {
179 fn new(use_large: bool) -> Self {
180 if use_large {
181 Self::Large(LargeStringBuilder::new())
182 } else {
183 Self::Regular(StringBuilder::new())
184 }
185 }
186 fn append_value(&mut self, v: &str) {
187 match self {
188 Self::Regular(b) => b.append_value(v),
189 Self::Large(b) => b.append_value(v),
190 }
191 }
192 fn append_null(&mut self) {
193 match self {
194 Self::Regular(b) => b.append_null(),
195 Self::Large(b) => b.append_null(),
196 }
197 }
198 fn append_default(&mut self) {
199 self.append_value("");
200 }
201 fn finish(&mut self) -> Arc<dyn Array> {
202 match self {
203 Self::Regular(b) => Arc::new(std::mem::take(b).finish()),
204 Self::Large(b) => Arc::new(std::mem::take(b).finish()),
205 }
206 }
207 fn len(&self) -> usize {
208 match self {
209 Self::Regular(b) => ArrayBuilder::len(b),
210 Self::Large(b) => ArrayBuilder::len(b),
211 }
212 }
213}
214
215enum BinaryBuilderInner {
216 Regular(BinaryBuilder),
217 Large(LargeBinaryBuilder),
218}
219
220impl BinaryBuilderInner {
221 fn new(use_large: bool) -> Self {
222 if use_large {
223 Self::Large(LargeBinaryBuilder::new())
224 } else {
225 Self::Regular(BinaryBuilder::new())
226 }
227 }
228 fn append_value(&mut self, v: &[u8]) {
229 match self {
230 Self::Regular(b) => b.append_value(v),
231 Self::Large(b) => b.append_value(v),
232 }
233 }
234 fn append_null(&mut self) {
235 match self {
236 Self::Regular(b) => b.append_null(),
237 Self::Large(b) => b.append_null(),
238 }
239 }
240 fn append_default(&mut self) {
241 self.append_value(b"");
242 }
243 fn finish(&mut self) -> Arc<dyn Array> {
244 match self {
245 Self::Regular(b) => Arc::new(std::mem::take(b).finish()),
246 Self::Large(b) => Arc::new(std::mem::take(b).finish()),
247 }
248 }
249 fn len(&self) -> usize {
250 match self {
251 Self::Regular(b) => ArrayBuilder::len(b),
252 Self::Large(b) => ArrayBuilder::len(b),
253 }
254 }
255}
256
257enum ListOffsets {
262 Regular(Vec<i32>),
263 Large(Vec<i64>),
264}
265
266impl ListOffsets {
267 fn new(use_large: bool) -> Self {
268 if use_large {
269 Self::Large(vec![0])
270 } else {
271 Self::Regular(vec![0])
272 }
273 }
274 fn push(&mut self, value: usize) {
275 match self {
276 Self::Regular(v) => v.push(value as i32),
277 Self::Large(v) => v.push(value as i64),
278 }
279 }
280 fn finish(self, values: Arc<dyn Array>, name: &str, nullable: bool) -> Arc<dyn Array> {
281 let field = Arc::new(Field::new(name, values.data_type().clone(), nullable));
282 match self {
283 Self::Regular(offsets) => {
284 let buf = Buffer::from_vec(offsets);
285 let data = ArrayData::builder(DataType::List(field))
286 .len(buf.len() / 4 - 1)
287 .add_buffer(buf)
288 .add_child_data(values.to_data())
289 .build()
290 .unwrap();
291 Arc::new(arrow_array::ListArray::from(data))
292 }
293 Self::Large(offsets) => {
294 let buf = Buffer::from_vec(offsets);
295 let data = ArrayData::builder(DataType::LargeList(field))
296 .len(buf.len() / 8 - 1)
297 .add_buffer(buf)
298 .add_child_data(values.to_data())
299 .build()
300 .unwrap();
301 Arc::new(arrow_array::LargeListArray::from(data))
302 }
303 }
304 }
305}
306
307enum RepeatedInner {
312 Int32 {
313 values_builder: PrimitiveBuilder<Int32Type>,
314 },
315 Int64 {
316 values_builder: PrimitiveBuilder<Int64Type>,
317 },
318 UInt32 {
319 values_builder: PrimitiveBuilder<UInt32Type>,
320 },
321 UInt64 {
322 values_builder: PrimitiveBuilder<UInt64Type>,
323 },
324 Float {
325 values_builder: PrimitiveBuilder<Float32Type>,
326 },
327 Double {
328 values_builder: PrimitiveBuilder<Float64Type>,
329 },
330 Bool {
331 values_builder: BooleanBuilder,
332 },
333 String {
334 values_builder: StringBuilderInner,
335 },
336 Bytes {
337 values_builder: BinaryBuilderInner,
338 },
339 Sint32 {
340 values_builder: PrimitiveBuilder<Int32Type>,
341 },
342 Sint64 {
343 values_builder: PrimitiveBuilder<Int64Type>,
344 },
345 Sfixed32 {
346 values_builder: PrimitiveBuilder<Int32Type>,
347 },
348 Sfixed64 {
349 values_builder: PrimitiveBuilder<Int64Type>,
350 },
351 Fixed32 {
352 values_builder: PrimitiveBuilder<UInt32Type>,
353 },
354 Fixed64 {
355 values_builder: PrimitiveBuilder<UInt64Type>,
356 },
357 EnumInt32 {
358 values_builder: PrimitiveBuilder<Int32Type>,
359 },
360 EnumString {
361 values_builder: StringBuilderInner,
362 enum_descriptor: EnumDescriptor,
363 },
364 EnumBinary {
365 values_builder: BinaryBuilderInner,
366 enum_descriptor: EnumDescriptor,
367 },
368 Message {
369 sub_decoder: MessageDecoder,
370 },
371 Timestamp {
372 values_builder: PrimitiveBuilder<Int64Type>,
373 unit: TimeUnit,
374 tz: Option<Arc<str>>,
375 },
376 Duration {
377 values_builder: PrimitiveBuilder<Int64Type>,
378 unit: TimeUnit,
379 },
380 Date {
381 values_builder: PrimitiveBuilder<Date32Type>,
382 },
383 TimeOfDay {
384 values_builder: PrimitiveBuilder<Int64Type>,
385 unit: TimeUnit,
386 },
387 WrapperDouble {
388 values_builder: PrimitiveBuilder<Float64Type>,
389 },
390 WrapperFloat {
391 values_builder: PrimitiveBuilder<Float32Type>,
392 },
393 WrapperInt64 {
394 values_builder: PrimitiveBuilder<Int64Type>,
395 },
396 WrapperUInt64 {
397 values_builder: PrimitiveBuilder<UInt64Type>,
398 },
399 WrapperInt32 {
400 values_builder: PrimitiveBuilder<Int32Type>,
401 },
402 WrapperUInt32 {
403 values_builder: PrimitiveBuilder<UInt32Type>,
404 },
405 WrapperBool {
406 values_builder: BooleanBuilder,
407 },
408 WrapperString {
409 values_builder: StringBuilderInner,
410 },
411 WrapperBytes {
412 values_builder: BinaryBuilderInner,
413 },
414}
415
416impl RepeatedInner {
417 fn decode(&mut self, wire_type: u8, buf: &[u8]) -> Result<usize, prost::DecodeError> {
418 match self {
419 Self::Int32 { values_builder, .. } => {
420 decode_repeated_varint(wire_type, buf, values_builder, |v| v as i32)
421 }
422 Self::Int64 { values_builder, .. } => {
423 decode_repeated_varint(wire_type, buf, values_builder, |v| v as i64)
424 }
425 Self::UInt32 { values_builder, .. } => {
426 decode_repeated_varint(wire_type, buf, values_builder, |v| v as u32)
427 }
428 Self::UInt64 { values_builder, .. } => {
429 decode_repeated_varint(wire_type, buf, values_builder, |v| v)
430 }
431 Self::EnumInt32 { values_builder, .. } => {
432 decode_repeated_varint(wire_type, buf, values_builder, |v| v as i32)
433 }
434 Self::Sint32 { values_builder, .. } => {
435 decode_repeated_varint(wire_type, buf, values_builder, decode_zigzag32)
436 }
437 Self::Sint64 { values_builder, .. } => {
438 decode_repeated_varint(wire_type, buf, values_builder, decode_zigzag64)
439 }
440 Self::Sfixed32 { values_builder, .. } => decode_repeated_fixed::<Int32Type, 4>(
441 wire_type,
442 5,
443 buf,
444 values_builder,
445 i32::from_le_bytes,
446 ),
447 Self::Sfixed64 { values_builder, .. } => decode_repeated_fixed::<Int64Type, 8>(
448 wire_type,
449 1,
450 buf,
451 values_builder,
452 i64::from_le_bytes,
453 ),
454 Self::Fixed32 { values_builder, .. } => decode_repeated_fixed::<UInt32Type, 4>(
455 wire_type,
456 5,
457 buf,
458 values_builder,
459 u32::from_le_bytes,
460 ),
461 Self::Fixed64 { values_builder, .. } => decode_repeated_fixed::<UInt64Type, 8>(
462 wire_type,
463 1,
464 buf,
465 values_builder,
466 u64::from_le_bytes,
467 ),
468 Self::Float { values_builder, .. } => decode_repeated_fixed::<Float32Type, 4>(
469 wire_type,
470 5,
471 buf,
472 values_builder,
473 f32::from_le_bytes,
474 ),
475 Self::Double { values_builder, .. } => decode_repeated_fixed::<Float64Type, 8>(
476 wire_type,
477 1,
478 buf,
479 values_builder,
480 f64::from_le_bytes,
481 ),
482 Self::Bool { values_builder, .. } => {
483 if wire_type == 2 {
484 let (data, total) = read_length_delimited(buf)?;
485 let mut p = 0;
486 while p < data.len() {
487 let (v, n) = decode_varint(&data[p..])?;
488 values_builder.append_value(v != 0);
489 p += n;
490 }
491 Ok(total)
492 } else if wire_type == 0 {
493 let (v, n) = decode_varint(buf)?;
494 values_builder.append_value(v != 0);
495 Ok(n)
496 } else {
497 skip_field(wire_type, buf)
498 }
499 }
500 Self::String { values_builder, .. } => {
501 if wire_type != 2 {
502 return skip_field(wire_type, buf);
503 }
504 let (data, total) = read_length_delimited(buf)?;
505 let s = std::str::from_utf8(data).map_err(|_| decode_error("invalid UTF-8"))?;
506 values_builder.append_value(s);
507 Ok(total)
508 }
509 Self::Bytes { values_builder, .. } => {
510 if wire_type != 2 {
511 return skip_field(wire_type, buf);
512 }
513 let (data, total) = read_length_delimited(buf)?;
514 values_builder.append_value(data);
515 Ok(total)
516 }
517 Self::EnumString {
518 values_builder,
519 enum_descriptor,
520 ..
521 } => {
522 if wire_type == 2 {
523 let (data, total) = read_length_delimited(buf)?;
524 let mut p = 0;
525 while p < data.len() {
526 let (v, n) = decode_varint(&data[p..])?;
527 values_builder.append_value(&enum_name(enum_descriptor, v as i32));
528 p += n;
529 }
530 Ok(total)
531 } else if wire_type == 0 {
532 let (v, n) = decode_varint(buf)?;
533 values_builder.append_value(&enum_name(enum_descriptor, v as i32));
534 Ok(n)
535 } else {
536 skip_field(wire_type, buf)
537 }
538 }
539 Self::EnumBinary {
540 values_builder,
541 enum_descriptor,
542 ..
543 } => {
544 if wire_type == 2 {
545 let (data, total) = read_length_delimited(buf)?;
546 let mut p = 0;
547 while p < data.len() {
548 let (v, n) = decode_varint(&data[p..])?;
549 values_builder
550 .append_value(enum_name(enum_descriptor, v as i32).as_bytes());
551 p += n;
552 }
553 Ok(total)
554 } else if wire_type == 0 {
555 let (v, n) = decode_varint(buf)?;
556 values_builder.append_value(enum_name(enum_descriptor, v as i32).as_bytes());
557 Ok(n)
558 } else {
559 skip_field(wire_type, buf)
560 }
561 }
562 Self::Message { sub_decoder, .. } => {
563 if wire_type != 2 {
564 return skip_field(wire_type, buf);
565 }
566 let (data, total) = read_length_delimited(buf)?;
567 sub_decoder.decode_row(data)?;
568 Ok(total)
569 }
570 Self::Timestamp {
571 values_builder,
572 unit,
573 ..
574 } => {
575 if wire_type != 2 {
576 return skip_field(wire_type, buf);
577 }
578 let (data, total) = read_length_delimited(buf)?;
579 let vals = decode_wkt_submessage(data, 2)?;
580 values_builder.append_value(convert_seconds_nanos_to_unit(
581 vals[0],
582 vals[1] as i32,
583 *unit,
584 "Timestamp",
585 ));
586 Ok(total)
587 }
588 Self::Duration {
589 values_builder,
590 unit,
591 ..
592 } => {
593 if wire_type != 2 {
594 return skip_field(wire_type, buf);
595 }
596 let (data, total) = read_length_delimited(buf)?;
597 let vals = decode_wkt_submessage(data, 2)?;
598 values_builder.append_value(convert_seconds_nanos_to_unit(
599 vals[0],
600 vals[1] as i32,
601 *unit,
602 "Duration",
603 ));
604 Ok(total)
605 }
606 Self::Date { values_builder, .. } => {
607 if wire_type != 2 {
608 return skip_field(wire_type, buf);
609 }
610 let (data, total) = read_length_delimited(buf)?;
611 let vals = decode_wkt_submessage(data, 3)?;
612 let (y, m, d) = (vals[0] as i32, vals[1] as i32, vals[2] as i32);
613 if y == 0 && m == 0 && d == 0 {
614 values_builder.append_value(0);
615 } else {
616 values_builder.append_value(
617 chrono::NaiveDate::from_ymd_opt(y, m as u32, d as u32)
618 .unwrap()
619 .num_days_from_ce()
620 - CE_OFFSET,
621 );
622 }
623 Ok(total)
624 }
625 Self::TimeOfDay {
626 values_builder,
627 unit,
628 ..
629 } => {
630 if wire_type != 2 {
631 return skip_field(wire_type, buf);
632 }
633 let (data, total) = read_length_delimited(buf)?;
634 let vals = decode_wkt_submessage(data, 4)?;
635 let total_seconds = vals[0] * 3600 + vals[1] * 60 + vals[2];
636 values_builder.append_value(convert_seconds_nanos_to_unit(
637 total_seconds,
638 vals[3] as i32,
639 *unit,
640 "TimeOfDay",
641 ));
642 Ok(total)
643 }
644 Self::WrapperDouble { values_builder, .. } => {
645 decode_repeated_wrapper_fixed64(wire_type, buf, values_builder, f64::from_le_bytes)
646 }
647 Self::WrapperFloat { values_builder, .. } => {
648 decode_repeated_wrapper_fixed32(wire_type, buf, values_builder, f32::from_le_bytes)
649 }
650 Self::WrapperInt64 { values_builder, .. } => {
651 decode_repeated_wrapper_varint(wire_type, buf, values_builder, |v| v as i64)
652 }
653 Self::WrapperUInt64 { values_builder, .. } => {
654 decode_repeated_wrapper_varint(wire_type, buf, values_builder, |v| v)
655 }
656 Self::WrapperInt32 { values_builder, .. } => {
657 decode_repeated_wrapper_varint(wire_type, buf, values_builder, |v| v as i32)
658 }
659 Self::WrapperUInt32 { values_builder, .. } => {
660 decode_repeated_wrapper_varint(wire_type, buf, values_builder, |v| v as u32)
661 }
662 Self::WrapperBool { values_builder, .. } => {
663 if wire_type != 2 {
664 return skip_field(wire_type, buf);
665 }
666 let (data, total) = read_length_delimited(buf)?;
667 let (v, _) = decode_wrapper_varint(data)?;
668 values_builder.append_value(v != 0);
669 Ok(total)
670 }
671 Self::WrapperString { values_builder, .. } => {
672 if wire_type != 2 {
673 return skip_field(wire_type, buf);
674 }
675 let (data, total) = read_length_delimited(buf)?;
676 let (v, found) = decode_wrapper_string(data)?;
677 if found {
678 values_builder.append_value(unsafe { std::str::from_utf8_unchecked(&v) });
679 } else {
680 values_builder.append_value("");
681 }
682 Ok(total)
683 }
684 Self::WrapperBytes { values_builder, .. } => {
685 if wire_type != 2 {
686 return skip_field(wire_type, buf);
687 }
688 let (data, total) = read_length_delimited(buf)?;
689 let (v, found) = decode_wrapper_bytes(data)?;
690 if found {
691 values_builder.append_value(&v);
692 } else {
693 values_builder.append_value(b"");
694 }
695 Ok(total)
696 }
697 }
698 }
699
700 fn len(&self) -> usize {
701 match self {
702 Self::Int32 { values_builder, .. }
703 | Self::Sint32 { values_builder, .. }
704 | Self::Sfixed32 { values_builder, .. }
705 | Self::EnumInt32 { values_builder, .. }
706 | Self::WrapperInt32 { values_builder, .. } => values_builder.len(),
707 Self::Int64 { values_builder, .. }
708 | Self::Sint64 { values_builder, .. }
709 | Self::Sfixed64 { values_builder, .. }
710 | Self::Timestamp { values_builder, .. }
711 | Self::Duration { values_builder, .. }
712 | Self::WrapperInt64 { values_builder, .. }
713 | Self::TimeOfDay { values_builder, .. } => values_builder.len(),
714 Self::UInt32 { values_builder, .. }
715 | Self::Fixed32 { values_builder, .. }
716 | Self::WrapperUInt32 { values_builder, .. } => values_builder.len(),
717 Self::UInt64 { values_builder, .. }
718 | Self::Fixed64 { values_builder, .. }
719 | Self::WrapperUInt64 { values_builder, .. } => values_builder.len(),
720 Self::Float { values_builder, .. } | Self::WrapperFloat { values_builder, .. } => {
721 values_builder.len()
722 }
723 Self::Double { values_builder, .. } | Self::WrapperDouble { values_builder, .. } => {
724 values_builder.len()
725 }
726 Self::Bool { values_builder, .. } | Self::WrapperBool { values_builder, .. } => {
727 values_builder.len()
728 }
729 Self::String { values_builder, .. }
730 | Self::EnumString { values_builder, .. }
731 | Self::WrapperString { values_builder, .. } => values_builder.len(),
732 Self::Bytes { values_builder, .. }
733 | Self::EnumBinary { values_builder, .. }
734 | Self::WrapperBytes { values_builder, .. } => values_builder.len(),
735 Self::Message { sub_decoder, .. } => sub_decoder.row_count(),
736 Self::Date { values_builder, .. } => values_builder.len(),
737 }
738 }
739
740 fn finish(&mut self) -> Arc<dyn Array> {
741 match self {
742 Self::Int32 { values_builder, .. }
743 | Self::Sint32 { values_builder, .. }
744 | Self::Sfixed32 { values_builder, .. }
745 | Self::EnumInt32 { values_builder, .. }
746 | Self::WrapperInt32 { values_builder, .. } => {
747 Arc::new(std::mem::take(values_builder).finish())
748 }
749 Self::Int64 { values_builder, .. }
750 | Self::Sint64 { values_builder, .. }
751 | Self::Sfixed64 { values_builder, .. }
752 | Self::WrapperInt64 { values_builder, .. } => {
753 Arc::new(std::mem::take(values_builder).finish())
754 }
755 Self::UInt32 { values_builder, .. }
756 | Self::Fixed32 { values_builder, .. }
757 | Self::WrapperUInt32 { values_builder, .. } => {
758 Arc::new(std::mem::take(values_builder).finish())
759 }
760 Self::UInt64 { values_builder, .. }
761 | Self::Fixed64 { values_builder, .. }
762 | Self::WrapperUInt64 { values_builder, .. } => {
763 Arc::new(std::mem::take(values_builder).finish())
764 }
765 Self::Float { values_builder, .. } | Self::WrapperFloat { values_builder, .. } => {
766 Arc::new(std::mem::take(values_builder).finish())
767 }
768 Self::Double { values_builder, .. } | Self::WrapperDouble { values_builder, .. } => {
769 Arc::new(std::mem::take(values_builder).finish())
770 }
771 Self::Bool { values_builder, .. } | Self::WrapperBool { values_builder, .. } => {
772 Arc::new(std::mem::take(values_builder).finish())
773 }
774 Self::String { values_builder, .. }
775 | Self::EnumString { values_builder, .. }
776 | Self::WrapperString { values_builder, .. } => values_builder.finish(),
777 Self::Bytes { values_builder, .. }
778 | Self::EnumBinary { values_builder, .. }
779 | Self::WrapperBytes { values_builder, .. } => values_builder.finish(),
780 Self::Message { sub_decoder, .. } => Arc::new(sub_decoder.build_struct_array(None)),
781 Self::Timestamp {
782 values_builder,
783 unit,
784 tz,
785 } => finish_timestamp(values_builder, *unit, tz),
786 Self::Duration {
787 values_builder,
788 unit,
789 } => finish_duration(values_builder, *unit),
790 Self::Date { values_builder, .. } => Arc::new(std::mem::take(values_builder).finish()),
791 Self::TimeOfDay {
792 values_builder,
793 unit,
794 ..
795 } => finish_time_of_day(values_builder, *unit),
796 }
797 }
798}
799
800enum FieldDecoder {
805 Int32 {
807 value: i32,
808 has_value: bool,
809 has_presence: bool,
810 builder: PrimitiveBuilder<Int32Type>,
811 },
812 Int64 {
813 value: i64,
814 has_value: bool,
815 has_presence: bool,
816 builder: PrimitiveBuilder<Int64Type>,
817 },
818 UInt32 {
819 value: u32,
820 has_value: bool,
821 has_presence: bool,
822 builder: PrimitiveBuilder<UInt32Type>,
823 },
824 UInt64 {
825 value: u64,
826 has_value: bool,
827 has_presence: bool,
828 builder: PrimitiveBuilder<UInt64Type>,
829 },
830 Sint32 {
831 value: i32,
832 has_value: bool,
833 has_presence: bool,
834 builder: PrimitiveBuilder<Int32Type>,
835 },
836 Sint64 {
837 value: i64,
838 has_value: bool,
839 has_presence: bool,
840 builder: PrimitiveBuilder<Int64Type>,
841 },
842 Sfixed32 {
843 value: i32,
844 has_value: bool,
845 has_presence: bool,
846 builder: PrimitiveBuilder<Int32Type>,
847 },
848 Sfixed64 {
849 value: i64,
850 has_value: bool,
851 has_presence: bool,
852 builder: PrimitiveBuilder<Int64Type>,
853 },
854 Fixed32 {
855 value: u32,
856 has_value: bool,
857 has_presence: bool,
858 builder: PrimitiveBuilder<UInt32Type>,
859 },
860 Fixed64 {
861 value: u64,
862 has_value: bool,
863 has_presence: bool,
864 builder: PrimitiveBuilder<UInt64Type>,
865 },
866 Float {
867 value: f32,
868 has_value: bool,
869 has_presence: bool,
870 builder: PrimitiveBuilder<Float32Type>,
871 },
872 Double {
873 value: f64,
874 has_value: bool,
875 has_presence: bool,
876 builder: PrimitiveBuilder<Float64Type>,
877 },
878 Bool {
879 value: bool,
880 has_value: bool,
881 has_presence: bool,
882 builder: BooleanBuilder,
883 },
884 String {
885 value: Vec<u8>,
886 has_value: bool,
887 has_presence: bool,
888 builder: StringBuilderInner,
889 },
890 Bytes {
891 value: Vec<u8>,
892 has_value: bool,
893 has_presence: bool,
894 builder: BinaryBuilderInner,
895 },
896 EnumInt32 {
897 value: i32,
898 has_value: bool,
899 has_presence: bool,
900 builder: PrimitiveBuilder<Int32Type>,
901 },
902 EnumString {
903 value: i32,
904 has_value: bool,
905 has_presence: bool,
906 builder: StringBuilderInner,
907 enum_descriptor: EnumDescriptor,
908 },
909 EnumBinary {
910 value: i32,
911 has_value: bool,
912 has_presence: bool,
913 builder: BinaryBuilderInner,
914 enum_descriptor: EnumDescriptor,
915 },
916
917 Timestamp {
919 seconds: i64,
920 nanos: i32,
921 has_value: bool,
922 builder: PrimitiveBuilder<Int64Type>,
923 unit: TimeUnit,
924 tz: Option<Arc<str>>,
925 },
926 Duration {
927 seconds: i64,
928 nanos: i32,
929 has_value: bool,
930 builder: PrimitiveBuilder<Int64Type>,
931 unit: TimeUnit,
932 },
933 Date {
934 year: i32,
935 month: i32,
936 day: i32,
937 has_value: bool,
938 builder: PrimitiveBuilder<Date32Type>,
939 },
940 TimeOfDay {
941 hours: i32,
942 minutes: i32,
943 seconds_val: i32,
944 nanos: i32,
945 has_value: bool,
946 builder: PrimitiveBuilder<Int64Type>,
947 unit: TimeUnit,
948 },
949
950 WrapperDouble {
952 value: f64,
953 has_value: bool,
954 builder: PrimitiveBuilder<Float64Type>,
955 },
956 WrapperFloat {
957 value: f32,
958 has_value: bool,
959 builder: PrimitiveBuilder<Float32Type>,
960 },
961 WrapperInt64 {
962 value: i64,
963 has_value: bool,
964 builder: PrimitiveBuilder<Int64Type>,
965 },
966 WrapperUInt64 {
967 value: u64,
968 has_value: bool,
969 builder: PrimitiveBuilder<UInt64Type>,
970 },
971 WrapperInt32 {
972 value: i32,
973 has_value: bool,
974 builder: PrimitiveBuilder<Int32Type>,
975 },
976 WrapperUInt32 {
977 value: u32,
978 has_value: bool,
979 builder: PrimitiveBuilder<UInt32Type>,
980 },
981 WrapperBool {
982 value: bool,
983 has_value: bool,
984 builder: BooleanBuilder,
985 },
986 WrapperString {
987 value: Vec<u8>,
988 has_value: bool,
989 builder: StringBuilderInner,
990 },
991 WrapperBytes {
992 value: Vec<u8>,
993 has_value: bool,
994 builder: BinaryBuilderInner,
995 },
996
997 Message {
999 sub_decoder: MessageDecoder,
1000 has_value: bool,
1001 is_valid: BooleanBuilder,
1002 },
1003
1004 Repeated {
1006 inner: RepeatedInner,
1007 offsets: ListOffsets,
1008 list_name: Arc<str>,
1009 list_nullable: bool,
1010 },
1011
1012 Map {
1014 key_decoder: Box<FieldDecoder>,
1015 value_decoder: Box<FieldDecoder>,
1016 offsets: Vec<i32>,
1017 map_value_name: Arc<str>,
1018 map_value_nullable: bool,
1019 },
1020}
1021
1022fn decode_wkt_submessage(buf: &[u8], max_field: u32) -> Result<[i64; 4], prost::DecodeError> {
1029 let mut vals = [0i64; 4];
1030 let mut pos = 0;
1031 while pos < buf.len() {
1032 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1033 pos += n;
1034 if fnum >= 1 && fnum <= max_field && wt == 0 {
1035 let (v, n) = decode_varint(&buf[pos..])?;
1036 vals[(fnum - 1) as usize] = v as i64;
1037 pos += n;
1038 } else {
1039 pos += skip_field(wt, &buf[pos..])?;
1040 }
1041 }
1042 Ok(vals)
1043}
1044
1045fn decode_wrapper_varint(buf: &[u8]) -> Result<(u64, bool), prost::DecodeError> {
1048 let mut val = 0u64;
1049 let mut found = false;
1050 let mut pos = 0;
1051 while pos < buf.len() {
1052 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1053 pos += n;
1054 if fnum == 1 && wt == 0 {
1055 let (v, n) = decode_varint(&buf[pos..])?;
1056 val = v;
1057 found = true;
1058 pos += n;
1059 } else {
1060 pos += skip_field(wt, &buf[pos..])?;
1061 }
1062 }
1063 Ok((val, found))
1064}
1065
1066fn decode_wrapper_fixed32(buf: &[u8]) -> Result<([u8; 4], bool), prost::DecodeError> {
1067 let mut val = [0u8; 4];
1068 let mut found = false;
1069 let mut pos = 0;
1070 while pos < buf.len() {
1071 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1072 pos += n;
1073 if fnum == 1 && wt == 5 {
1074 if buf.len() < pos + 4 {
1075 return Err(decode_error("unexpected EOF"));
1076 }
1077 val.copy_from_slice(&buf[pos..pos + 4]);
1078 found = true;
1079 pos += 4;
1080 } else {
1081 pos += skip_field(wt, &buf[pos..])?;
1082 }
1083 }
1084 Ok((val, found))
1085}
1086
1087fn decode_wrapper_fixed64(buf: &[u8]) -> Result<([u8; 8], bool), prost::DecodeError> {
1088 let mut val = [0u8; 8];
1089 let mut found = false;
1090 let mut pos = 0;
1091 while pos < buf.len() {
1092 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1093 pos += n;
1094 if fnum == 1 && wt == 1 {
1095 if buf.len() < pos + 8 {
1096 return Err(decode_error("unexpected EOF"));
1097 }
1098 val.copy_from_slice(&buf[pos..pos + 8]);
1099 found = true;
1100 pos += 8;
1101 } else {
1102 pos += skip_field(wt, &buf[pos..])?;
1103 }
1104 }
1105 Ok((val, found))
1106}
1107
1108fn decode_wrapper_string(buf: &[u8]) -> Result<(Vec<u8>, bool), prost::DecodeError> {
1109 let mut val = Vec::new();
1110 let mut found = false;
1111 let mut pos = 0;
1112 while pos < buf.len() {
1113 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1114 pos += n;
1115 if fnum == 1 && wt == 2 {
1116 let (data, consumed) = read_length_delimited(&buf[pos..])?;
1117 std::str::from_utf8(data).map_err(|_| decode_error("invalid UTF-8"))?;
1118 val.clear();
1119 val.extend_from_slice(data);
1120 found = true;
1121 pos += consumed;
1122 } else {
1123 pos += skip_field(wt, &buf[pos..])?;
1124 }
1125 }
1126 Ok((val, found))
1127}
1128
1129fn decode_wrapper_bytes(buf: &[u8]) -> Result<(Vec<u8>, bool), prost::DecodeError> {
1130 let mut val = Vec::new();
1131 let mut found = false;
1132 let mut pos = 0;
1133 while pos < buf.len() {
1134 let (fnum, wt, n) = decode_tag(&buf[pos..])?;
1135 pos += n;
1136 if fnum == 1 && wt == 2 {
1137 let (data, consumed) = read_length_delimited(&buf[pos..])?;
1138 val.clear();
1139 val.extend_from_slice(data);
1140 found = true;
1141 pos += consumed;
1142 } else {
1143 pos += skip_field(wt, &buf[pos..])?;
1144 }
1145 }
1146 Ok((val, found))
1147}
1148
1149macro_rules! flush_primitive {
1154 ($value:expr, $has_value:expr, $has_presence:expr, $builder:expr, $default:expr) => {
1155 if *$has_value {
1156 $builder.append_value(*$value);
1157 } else if *$has_presence {
1158 $builder.append_null();
1159 } else {
1160 $builder.append_value($default);
1161 }
1162 *$has_value = false;
1163 *$value = $default;
1164 };
1165}
1166
1167fn finish_primitive<T: ArrowPrimitiveType>(builder: &mut PrimitiveBuilder<T>) -> Arc<dyn Array> {
1168 Arc::new(std::mem::take(builder).finish())
1169}
1170
1171fn finish_timestamp(
1172 builder: &mut PrimitiveBuilder<Int64Type>,
1173 unit: TimeUnit,
1174 tz: &Option<Arc<str>>,
1175) -> Arc<dyn Array> {
1176 let values = std::mem::take(builder).finish();
1177 let dt = DataType::Timestamp(unit, tz.clone());
1178 let data = ArrayData::builder(dt)
1179 .len(values.len())
1180 .add_buffer(values.values().inner().clone())
1181 .null_bit_buffer(values.nulls().map(|n| n.buffer().clone()))
1182 .build()
1183 .unwrap();
1184 arrow_array::make_array(data)
1185}
1186
1187fn finish_duration(builder: &mut PrimitiveBuilder<Int64Type>, unit: TimeUnit) -> Arc<dyn Array> {
1188 let values = std::mem::take(builder).finish();
1189 let dt = DataType::Duration(unit);
1190 let data = ArrayData::builder(dt)
1191 .len(values.len())
1192 .add_buffer(values.values().inner().clone())
1193 .null_bit_buffer(values.nulls().map(|n| n.buffer().clone()))
1194 .build()
1195 .unwrap();
1196 arrow_array::make_array(data)
1197}
1198
1199fn finish_time_of_day(builder: &mut PrimitiveBuilder<Int64Type>, unit: TimeUnit) -> Arc<dyn Array> {
1200 let values = std::mem::take(builder).finish();
1201 let dt = match unit {
1202 TimeUnit::Second => DataType::Time32(TimeUnit::Second),
1203 TimeUnit::Millisecond => DataType::Time32(TimeUnit::Millisecond),
1204 TimeUnit::Microsecond => DataType::Time64(TimeUnit::Microsecond),
1205 TimeUnit::Nanosecond => DataType::Time64(TimeUnit::Nanosecond),
1206 };
1207 if matches!(unit, TimeUnit::Second | TimeUnit::Millisecond) {
1208 let i32_values: Vec<Option<i32>> = (0..values.len())
1209 .map(|i| {
1210 if values.is_null(i) {
1211 None
1212 } else {
1213 let v = values.value(i);
1214 Some(i32::try_from(v).unwrap_or(if v > 0 { i32::MAX } else { i32::MIN }))
1215 }
1216 })
1217 .collect();
1218 let i32_array = arrow_array::Int32Array::from(i32_values);
1219 let data = ArrayData::builder(dt)
1220 .len(i32_array.len())
1221 .add_buffer(i32_array.values().inner().clone())
1222 .null_bit_buffer(i32_array.nulls().map(|n| n.buffer().clone()))
1223 .build()
1224 .unwrap();
1225 arrow_array::make_array(data)
1226 } else {
1227 let data = ArrayData::builder(dt)
1228 .len(values.len())
1229 .add_buffer(values.values().inner().clone())
1230 .null_bit_buffer(values.nulls().map(|n| n.buffer().clone()))
1231 .build()
1232 .unwrap();
1233 arrow_array::make_array(data)
1234 }
1235}
1236
1237impl FieldDecoder {
1242 fn decode(&mut self, wire_type: u8, buf: &[u8]) -> Result<usize, prost::DecodeError> {
1243 match self {
1244 Self::Int32 {
1245 value, has_value, ..
1246 }
1247 | Self::EnumInt32 {
1248 value, has_value, ..
1249 } => {
1250 if wire_type != 0 {
1251 return skip_field(wire_type, buf);
1252 }
1253 let (v, n) = decode_varint(buf)?;
1254 *value = v as i32;
1255 *has_value = true;
1256 Ok(n)
1257 }
1258 Self::EnumString {
1259 value, has_value, ..
1260 }
1261 | Self::EnumBinary {
1262 value, has_value, ..
1263 } => {
1264 if wire_type != 0 {
1265 return skip_field(wire_type, buf);
1266 }
1267 let (v, n) = decode_varint(buf)?;
1268 *value = v as i32;
1269 *has_value = true;
1270 Ok(n)
1271 }
1272 Self::Int64 {
1273 value, has_value, ..
1274 } => {
1275 if wire_type != 0 {
1276 return skip_field(wire_type, buf);
1277 }
1278 let (v, n) = decode_varint(buf)?;
1279 *value = v as i64;
1280 *has_value = true;
1281 Ok(n)
1282 }
1283 Self::UInt32 {
1284 value, has_value, ..
1285 } => {
1286 if wire_type != 0 {
1287 return skip_field(wire_type, buf);
1288 }
1289 let (v, n) = decode_varint(buf)?;
1290 *value = v as u32;
1291 *has_value = true;
1292 Ok(n)
1293 }
1294 Self::UInt64 {
1295 value, has_value, ..
1296 } => {
1297 if wire_type != 0 {
1298 return skip_field(wire_type, buf);
1299 }
1300 let (v, n) = decode_varint(buf)?;
1301 *value = v;
1302 *has_value = true;
1303 Ok(n)
1304 }
1305 Self::Sint32 {
1306 value, has_value, ..
1307 } => {
1308 if wire_type != 0 {
1309 return skip_field(wire_type, buf);
1310 }
1311 let (v, n) = decode_varint(buf)?;
1312 *value = decode_zigzag32(v);
1313 *has_value = true;
1314 Ok(n)
1315 }
1316 Self::Sint64 {
1317 value, has_value, ..
1318 } => {
1319 if wire_type != 0 {
1320 return skip_field(wire_type, buf);
1321 }
1322 let (v, n) = decode_varint(buf)?;
1323 *value = decode_zigzag64(v);
1324 *has_value = true;
1325 Ok(n)
1326 }
1327 Self::Sfixed32 {
1328 value, has_value, ..
1329 } => {
1330 if wire_type != 5 {
1331 return skip_field(wire_type, buf);
1332 }
1333 if buf.len() < 4 {
1334 return Err(decode_error("unexpected EOF"));
1335 }
1336 *value = i32::from_le_bytes(buf[..4].try_into().unwrap());
1337 *has_value = true;
1338 Ok(4)
1339 }
1340 Self::Sfixed64 {
1341 value, has_value, ..
1342 } => {
1343 if wire_type != 1 {
1344 return skip_field(wire_type, buf);
1345 }
1346 if buf.len() < 8 {
1347 return Err(decode_error("unexpected EOF"));
1348 }
1349 *value = i64::from_le_bytes(buf[..8].try_into().unwrap());
1350 *has_value = true;
1351 Ok(8)
1352 }
1353 Self::Fixed32 {
1354 value, has_value, ..
1355 } => {
1356 if wire_type != 5 {
1357 return skip_field(wire_type, buf);
1358 }
1359 if buf.len() < 4 {
1360 return Err(decode_error("unexpected EOF"));
1361 }
1362 *value = u32::from_le_bytes(buf[..4].try_into().unwrap());
1363 *has_value = true;
1364 Ok(4)
1365 }
1366 Self::Fixed64 {
1367 value, has_value, ..
1368 } => {
1369 if wire_type != 1 {
1370 return skip_field(wire_type, buf);
1371 }
1372 if buf.len() < 8 {
1373 return Err(decode_error("unexpected EOF"));
1374 }
1375 *value = u64::from_le_bytes(buf[..8].try_into().unwrap());
1376 *has_value = true;
1377 Ok(8)
1378 }
1379 Self::Float {
1380 value, has_value, ..
1381 } => {
1382 if wire_type != 5 {
1383 return skip_field(wire_type, buf);
1384 }
1385 if buf.len() < 4 {
1386 return Err(decode_error("unexpected EOF"));
1387 }
1388 *value = f32::from_le_bytes(buf[..4].try_into().unwrap());
1389 *has_value = true;
1390 Ok(4)
1391 }
1392 Self::Double {
1393 value, has_value, ..
1394 } => {
1395 if wire_type != 1 {
1396 return skip_field(wire_type, buf);
1397 }
1398 if buf.len() < 8 {
1399 return Err(decode_error("unexpected EOF"));
1400 }
1401 *value = f64::from_le_bytes(buf[..8].try_into().unwrap());
1402 *has_value = true;
1403 Ok(8)
1404 }
1405 Self::Bool {
1406 value, has_value, ..
1407 } => {
1408 if wire_type != 0 {
1409 return skip_field(wire_type, buf);
1410 }
1411 let (v, n) = decode_varint(buf)?;
1412 *value = v != 0;
1413 *has_value = true;
1414 Ok(n)
1415 }
1416 Self::String {
1417 value, has_value, ..
1418 } => {
1419 if wire_type != 2 {
1420 return skip_field(wire_type, buf);
1421 }
1422 let (data, total) = read_length_delimited(buf)?;
1423 std::str::from_utf8(data).map_err(|_| decode_error("invalid UTF-8"))?;
1424 value.clear();
1425 value.extend_from_slice(data);
1426 *has_value = true;
1427 Ok(total)
1428 }
1429 Self::Bytes {
1430 value, has_value, ..
1431 } => {
1432 if wire_type != 2 {
1433 return skip_field(wire_type, buf);
1434 }
1435 let (data, total) = read_length_delimited(buf)?;
1436 value.clear();
1437 value.extend_from_slice(data);
1438 *has_value = true;
1439 Ok(total)
1440 }
1441 Self::Timestamp {
1443 seconds,
1444 nanos,
1445 has_value,
1446 ..
1447 } => {
1448 if wire_type != 2 {
1449 return skip_field(wire_type, buf);
1450 }
1451 let (data, total) = read_length_delimited(buf)?;
1452 let vals = decode_wkt_submessage(data, 2)?;
1453 *seconds = vals[0];
1454 *nanos = vals[1] as i32;
1455 *has_value = true;
1456 Ok(total)
1457 }
1458 Self::Duration {
1459 seconds,
1460 nanos,
1461 has_value,
1462 ..
1463 } => {
1464 if wire_type != 2 {
1465 return skip_field(wire_type, buf);
1466 }
1467 let (data, total) = read_length_delimited(buf)?;
1468 let vals = decode_wkt_submessage(data, 2)?;
1469 *seconds = vals[0];
1470 *nanos = vals[1] as i32;
1471 *has_value = true;
1472 Ok(total)
1473 }
1474 Self::Date {
1475 year,
1476 month,
1477 day,
1478 has_value,
1479 ..
1480 } => {
1481 if wire_type != 2 {
1482 return skip_field(wire_type, buf);
1483 }
1484 let (data, total) = read_length_delimited(buf)?;
1485 let vals = decode_wkt_submessage(data, 3)?;
1486 *year = vals[0] as i32;
1487 *month = vals[1] as i32;
1488 *day = vals[2] as i32;
1489 *has_value = true;
1490 Ok(total)
1491 }
1492 Self::TimeOfDay {
1493 hours,
1494 minutes,
1495 seconds_val,
1496 nanos,
1497 has_value,
1498 ..
1499 } => {
1500 if wire_type != 2 {
1501 return skip_field(wire_type, buf);
1502 }
1503 let (data, total) = read_length_delimited(buf)?;
1504 let vals = decode_wkt_submessage(data, 4)?;
1505 *hours = vals[0] as i32;
1506 *minutes = vals[1] as i32;
1507 *seconds_val = vals[2] as i32;
1508 *nanos = vals[3] as i32;
1509 *has_value = true;
1510 Ok(total)
1511 }
1512 Self::WrapperDouble {
1514 value, has_value, ..
1515 } => {
1516 if wire_type != 2 {
1517 return skip_field(wire_type, buf);
1518 }
1519 let (data, total) = read_length_delimited(buf)?;
1520 let (bytes, found) = decode_wrapper_fixed64(data)?;
1521 if found {
1522 *value = f64::from_le_bytes(bytes);
1523 }
1524 *has_value = true;
1525 Ok(total)
1526 }
1527 Self::WrapperFloat {
1528 value, has_value, ..
1529 } => {
1530 if wire_type != 2 {
1531 return skip_field(wire_type, buf);
1532 }
1533 let (data, total) = read_length_delimited(buf)?;
1534 let (bytes, found) = decode_wrapper_fixed32(data)?;
1535 if found {
1536 *value = f32::from_le_bytes(bytes);
1537 }
1538 *has_value = true;
1539 Ok(total)
1540 }
1541 Self::WrapperInt64 {
1542 value, has_value, ..
1543 } => {
1544 if wire_type != 2 {
1545 return skip_field(wire_type, buf);
1546 }
1547 let (data, total) = read_length_delimited(buf)?;
1548 let (v, _) = decode_wrapper_varint(data)?;
1549 *value = v as i64;
1550 *has_value = true;
1551 Ok(total)
1552 }
1553 Self::WrapperUInt64 {
1554 value, has_value, ..
1555 } => {
1556 if wire_type != 2 {
1557 return skip_field(wire_type, buf);
1558 }
1559 let (data, total) = read_length_delimited(buf)?;
1560 let (v, _) = decode_wrapper_varint(data)?;
1561 *value = v;
1562 *has_value = true;
1563 Ok(total)
1564 }
1565 Self::WrapperInt32 {
1566 value, has_value, ..
1567 } => {
1568 if wire_type != 2 {
1569 return skip_field(wire_type, buf);
1570 }
1571 let (data, total) = read_length_delimited(buf)?;
1572 let (v, _) = decode_wrapper_varint(data)?;
1573 *value = v as i32;
1574 *has_value = true;
1575 Ok(total)
1576 }
1577 Self::WrapperUInt32 {
1578 value, has_value, ..
1579 } => {
1580 if wire_type != 2 {
1581 return skip_field(wire_type, buf);
1582 }
1583 let (data, total) = read_length_delimited(buf)?;
1584 let (v, _) = decode_wrapper_varint(data)?;
1585 *value = v as u32;
1586 *has_value = true;
1587 Ok(total)
1588 }
1589 Self::WrapperBool {
1590 value, has_value, ..
1591 } => {
1592 if wire_type != 2 {
1593 return skip_field(wire_type, buf);
1594 }
1595 let (data, total) = read_length_delimited(buf)?;
1596 let (v, _) = decode_wrapper_varint(data)?;
1597 *value = v != 0;
1598 *has_value = true;
1599 Ok(total)
1600 }
1601 Self::WrapperString {
1602 value, has_value, ..
1603 } => {
1604 if wire_type != 2 {
1605 return skip_field(wire_type, buf);
1606 }
1607 let (data, total) = read_length_delimited(buf)?;
1608 let (v, _) = decode_wrapper_string(data)?;
1609 value.clear();
1610 value.extend_from_slice(&v);
1611 *has_value = true;
1612 Ok(total)
1613 }
1614 Self::WrapperBytes {
1615 value, has_value, ..
1616 } => {
1617 if wire_type != 2 {
1618 return skip_field(wire_type, buf);
1619 }
1620 let (data, total) = read_length_delimited(buf)?;
1621 let (v, _) = decode_wrapper_bytes(data)?;
1622 value.clear();
1623 value.extend_from_slice(&v);
1624 *has_value = true;
1625 Ok(total)
1626 }
1627 Self::Message {
1629 sub_decoder,
1630 has_value,
1631 ..
1632 } => {
1633 if wire_type != 2 {
1634 return skip_field(wire_type, buf);
1635 }
1636 let (data, total) = read_length_delimited(buf)?;
1637 if !*has_value {
1641 *has_value = true;
1642 }
1643 sub_decoder.decode_message_bytes(data)?;
1644 Ok(total)
1645 }
1646 Self::Repeated { inner, .. } => inner.decode(wire_type, buf),
1648 Self::Map {
1650 key_decoder,
1651 value_decoder,
1652 ..
1653 } => {
1654 if wire_type != 2 {
1655 return skip_field(wire_type, buf);
1656 }
1657 let (data, total) = read_length_delimited(buf)?;
1658 let mut pos = 0;
1660 while pos < data.len() {
1661 let (fnum, wt, n) = decode_tag(&data[pos..])?;
1662 pos += n;
1663 if fnum == 1 {
1664 pos += key_decoder.decode(wt, &data[pos..])?;
1665 } else if fnum == 2 {
1666 pos += value_decoder.decode(wt, &data[pos..])?;
1667 } else {
1668 pos += skip_field(wt, &data[pos..])?;
1669 }
1670 }
1671 key_decoder.flush();
1673 value_decoder.flush();
1674 Ok(total)
1675 }
1676 }
1677 }
1678
1679 fn flush(&mut self) {
1680 match self {
1681 Self::Int32 {
1682 value,
1683 has_value,
1684 has_presence,
1685 builder,
1686 }
1687 | Self::Sint32 {
1688 value,
1689 has_value,
1690 has_presence,
1691 builder,
1692 }
1693 | Self::Sfixed32 {
1694 value,
1695 has_value,
1696 has_presence,
1697 builder,
1698 }
1699 | Self::EnumInt32 {
1700 value,
1701 has_value,
1702 has_presence,
1703 builder,
1704 } => {
1705 flush_primitive!(value, has_value, has_presence, builder, 0i32);
1706 }
1707 Self::Int64 {
1708 value,
1709 has_value,
1710 has_presence,
1711 builder,
1712 }
1713 | Self::Sint64 {
1714 value,
1715 has_value,
1716 has_presence,
1717 builder,
1718 }
1719 | Self::Sfixed64 {
1720 value,
1721 has_value,
1722 has_presence,
1723 builder,
1724 } => {
1725 flush_primitive!(value, has_value, has_presence, builder, 0i64);
1726 }
1727 Self::UInt32 {
1728 value,
1729 has_value,
1730 has_presence,
1731 builder,
1732 }
1733 | Self::Fixed32 {
1734 value,
1735 has_value,
1736 has_presence,
1737 builder,
1738 } => {
1739 flush_primitive!(value, has_value, has_presence, builder, 0u32);
1740 }
1741 Self::UInt64 {
1742 value,
1743 has_value,
1744 has_presence,
1745 builder,
1746 }
1747 | Self::Fixed64 {
1748 value,
1749 has_value,
1750 has_presence,
1751 builder,
1752 } => {
1753 flush_primitive!(value, has_value, has_presence, builder, 0u64);
1754 }
1755 Self::Float {
1756 value,
1757 has_value,
1758 has_presence,
1759 builder,
1760 } => {
1761 flush_primitive!(value, has_value, has_presence, builder, 0.0f32);
1762 }
1763 Self::Double {
1764 value,
1765 has_value,
1766 has_presence,
1767 builder,
1768 } => {
1769 flush_primitive!(value, has_value, has_presence, builder, 0.0f64);
1770 }
1771 Self::EnumString {
1772 value,
1773 has_value,
1774 has_presence,
1775 builder,
1776 enum_descriptor,
1777 } => {
1778 if *has_value {
1779 builder.append_value(&enum_name(enum_descriptor, *value));
1780 } else if *has_presence {
1781 builder.append_null();
1782 } else {
1783 builder.append_value(&enum_name(enum_descriptor, 0));
1784 }
1785 *has_value = false;
1786 *value = 0;
1787 }
1788 Self::EnumBinary {
1789 value,
1790 has_value,
1791 has_presence,
1792 builder,
1793 enum_descriptor,
1794 } => {
1795 if *has_value {
1796 builder.append_value(enum_name(enum_descriptor, *value).as_bytes());
1797 } else if *has_presence {
1798 builder.append_null();
1799 } else {
1800 builder.append_value(enum_name(enum_descriptor, 0).as_bytes());
1801 }
1802 *has_value = false;
1803 *value = 0;
1804 }
1805 Self::Bool {
1806 value,
1807 has_value,
1808 has_presence,
1809 builder,
1810 } => {
1811 if *has_value {
1812 builder.append_value(*value);
1813 } else if *has_presence {
1814 builder.append_null();
1815 } else {
1816 builder.append_value(false);
1817 }
1818 *has_value = false;
1819 *value = false;
1820 }
1821 Self::String {
1822 value,
1823 has_value,
1824 has_presence,
1825 builder,
1826 } => {
1827 if *has_value {
1828 builder.append_value(unsafe { std::str::from_utf8_unchecked(value) });
1829 } else if *has_presence {
1830 builder.append_null();
1831 } else {
1832 builder.append_default();
1833 }
1834 *has_value = false;
1835 value.clear();
1836 }
1837 Self::Bytes {
1838 value,
1839 has_value,
1840 has_presence,
1841 builder,
1842 } => {
1843 if *has_value {
1844 builder.append_value(value.as_slice());
1845 } else if *has_presence {
1846 builder.append_null();
1847 } else {
1848 builder.append_default();
1849 }
1850 *has_value = false;
1851 value.clear();
1852 }
1853 Self::Timestamp {
1855 seconds,
1856 nanos,
1857 has_value,
1858 builder,
1859 unit,
1860 ..
1861 } => {
1862 if *has_value {
1863 builder.append_value(convert_seconds_nanos_to_unit(
1864 *seconds,
1865 *nanos,
1866 *unit,
1867 "Timestamp",
1868 ));
1869 } else {
1870 builder.append_null();
1871 }
1872 *has_value = false;
1873 *seconds = 0;
1874 *nanos = 0;
1875 }
1876 Self::Duration {
1877 seconds,
1878 nanos,
1879 has_value,
1880 builder,
1881 unit,
1882 ..
1883 } => {
1884 if *has_value {
1885 builder.append_value(convert_seconds_nanos_to_unit(
1886 *seconds, *nanos, *unit, "Duration",
1887 ));
1888 } else {
1889 builder.append_null();
1890 }
1891 *has_value = false;
1892 *seconds = 0;
1893 *nanos = 0;
1894 }
1895 Self::Date {
1896 year,
1897 month,
1898 day,
1899 has_value,
1900 builder,
1901 } => {
1902 if *has_value {
1903 if *year == 0 && *month == 0 && *day == 0 {
1904 builder.append_value(0);
1905 } else {
1906 builder.append_value(
1907 chrono::NaiveDate::from_ymd_opt(*year, *month as u32, *day as u32)
1908 .unwrap()
1909 .num_days_from_ce()
1910 - CE_OFFSET,
1911 );
1912 }
1913 } else {
1914 builder.append_null();
1915 }
1916 *has_value = false;
1917 *year = 0;
1918 *month = 0;
1919 *day = 0;
1920 }
1921 Self::TimeOfDay {
1922 hours,
1923 minutes,
1924 seconds_val,
1925 nanos,
1926 has_value,
1927 builder,
1928 unit,
1929 } => {
1930 if *has_value {
1931 let total_seconds = i64::from(*hours) * 3600
1932 + i64::from(*minutes) * 60
1933 + i64::from(*seconds_val);
1934 builder.append_value(convert_seconds_nanos_to_unit(
1935 total_seconds,
1936 *nanos,
1937 *unit,
1938 "TimeOfDay",
1939 ));
1940 } else {
1941 builder.append_null();
1942 }
1943 *has_value = false;
1944 *hours = 0;
1945 *minutes = 0;
1946 *seconds_val = 0;
1947 *nanos = 0;
1948 }
1949 Self::WrapperDouble {
1951 value,
1952 has_value,
1953 builder,
1954 } => {
1955 if *has_value {
1956 builder.append_value(*value);
1957 } else {
1958 builder.append_null();
1959 }
1960 *has_value = false;
1961 *value = 0.0;
1962 }
1963 Self::WrapperFloat {
1964 value,
1965 has_value,
1966 builder,
1967 } => {
1968 if *has_value {
1969 builder.append_value(*value);
1970 } else {
1971 builder.append_null();
1972 }
1973 *has_value = false;
1974 *value = 0.0;
1975 }
1976 Self::WrapperInt64 {
1977 value,
1978 has_value,
1979 builder,
1980 } => {
1981 if *has_value {
1982 builder.append_value(*value);
1983 } else {
1984 builder.append_null();
1985 }
1986 *has_value = false;
1987 *value = 0;
1988 }
1989 Self::WrapperUInt64 {
1990 value,
1991 has_value,
1992 builder,
1993 } => {
1994 if *has_value {
1995 builder.append_value(*value);
1996 } else {
1997 builder.append_null();
1998 }
1999 *has_value = false;
2000 *value = 0;
2001 }
2002 Self::WrapperInt32 {
2003 value,
2004 has_value,
2005 builder,
2006 } => {
2007 if *has_value {
2008 builder.append_value(*value);
2009 } else {
2010 builder.append_null();
2011 }
2012 *has_value = false;
2013 *value = 0;
2014 }
2015 Self::WrapperUInt32 {
2016 value,
2017 has_value,
2018 builder,
2019 } => {
2020 if *has_value {
2021 builder.append_value(*value);
2022 } else {
2023 builder.append_null();
2024 }
2025 *has_value = false;
2026 *value = 0;
2027 }
2028 Self::WrapperBool {
2029 value,
2030 has_value,
2031 builder,
2032 } => {
2033 if *has_value {
2034 builder.append_value(*value);
2035 } else {
2036 builder.append_null();
2037 }
2038 *has_value = false;
2039 *value = false;
2040 }
2041 Self::WrapperString {
2042 value,
2043 has_value,
2044 builder,
2045 } => {
2046 if *has_value {
2047 builder.append_value(unsafe { std::str::from_utf8_unchecked(value) });
2048 } else {
2049 builder.append_null();
2050 }
2051 *has_value = false;
2052 value.clear();
2053 }
2054 Self::WrapperBytes {
2055 value,
2056 has_value,
2057 builder,
2058 } => {
2059 if *has_value {
2060 builder.append_value(value.as_slice());
2061 } else {
2062 builder.append_null();
2063 }
2064 *has_value = false;
2065 value.clear();
2066 }
2067 Self::Message {
2069 sub_decoder,
2070 has_value,
2071 is_valid,
2072 } => {
2073 if *has_value {
2074 is_valid.append_value(true);
2075 sub_decoder.flush_row();
2076 } else {
2077 is_valid.append_value(false);
2078 sub_decoder.flush_defaults();
2079 }
2080 *has_value = false;
2081 }
2082 Self::Repeated { inner, offsets, .. } => {
2084 offsets.push(inner.len());
2085 }
2086 Self::Map {
2088 key_decoder,
2089 offsets,
2090 ..
2091 } => {
2092 let count = match key_decoder.as_ref() {
2093 FieldDecoder::Int32 { builder, .. } => ArrayBuilder::len(builder),
2094 FieldDecoder::Int64 { builder, .. } => ArrayBuilder::len(builder),
2095 FieldDecoder::UInt32 { builder, .. } => ArrayBuilder::len(builder),
2096 FieldDecoder::UInt64 { builder, .. } => ArrayBuilder::len(builder),
2097 FieldDecoder::Sint32 { builder, .. } => ArrayBuilder::len(builder),
2098 FieldDecoder::Sint64 { builder, .. } => ArrayBuilder::len(builder),
2099 FieldDecoder::Bool { builder, .. } => ArrayBuilder::len(builder),
2100 FieldDecoder::String { builder, .. } => builder.len(),
2101 _ => *offsets.last().unwrap() as usize,
2102 };
2103 offsets.push(count as i32);
2104 }
2105 }
2106 }
2107
2108 fn finish(&mut self, nullable: bool) -> (Field, Arc<dyn Array>) {
2109 let array: Arc<dyn Array> = match self {
2112 Self::Int32 { builder, .. } | Self::EnumInt32 { builder, .. } => {
2113 finish_primitive(builder)
2114 }
2115 Self::Int64 { builder, .. } => finish_primitive(builder),
2116 Self::UInt32 { builder, .. } => finish_primitive(builder),
2117 Self::UInt64 { builder, .. } => finish_primitive(builder),
2118 Self::Sint32 { builder, .. } | Self::Sfixed32 { builder, .. } => {
2119 finish_primitive(builder)
2120 }
2121 Self::Sint64 { builder, .. } | Self::Sfixed64 { builder, .. } => {
2122 finish_primitive(builder)
2123 }
2124 Self::Fixed32 { builder, .. } => finish_primitive(builder),
2125 Self::Fixed64 { builder, .. } => finish_primitive(builder),
2126 Self::Float { builder, .. } => finish_primitive(builder),
2127 Self::Double { builder, .. } => finish_primitive(builder),
2128 Self::Bool { builder, .. } => Arc::new(std::mem::take(builder).finish()),
2129 Self::String { builder, .. } | Self::EnumString { builder, .. } => builder.finish(),
2130 Self::Bytes { builder, .. } | Self::EnumBinary { builder, .. } => builder.finish(),
2131 Self::Timestamp {
2132 builder, unit, tz, ..
2133 } => finish_timestamp(builder, *unit, tz),
2134 Self::Duration { builder, unit, .. } => finish_duration(builder, *unit),
2135 Self::Date { builder, .. } => finish_primitive(builder),
2136 Self::TimeOfDay { builder, unit, .. } => finish_time_of_day(builder, *unit),
2137 Self::WrapperDouble { builder, .. } => finish_primitive(builder),
2138 Self::WrapperFloat { builder, .. } => finish_primitive(builder),
2139 Self::WrapperInt64 { builder, .. } => finish_primitive(builder),
2140 Self::WrapperUInt64 { builder, .. } => finish_primitive(builder),
2141 Self::WrapperInt32 { builder, .. } => finish_primitive(builder),
2142 Self::WrapperUInt32 { builder, .. } => finish_primitive(builder),
2143 Self::WrapperBool { builder, .. } => Arc::new(std::mem::take(builder).finish()),
2144 Self::WrapperString { builder, .. } => builder.finish(),
2145 Self::WrapperBytes { builder, .. } => builder.finish(),
2146 Self::Message {
2147 sub_decoder,
2148 is_valid,
2149 ..
2150 } => Arc::new(sub_decoder.build_struct_array(Some(std::mem::take(is_valid).finish()))),
2151 Self::Repeated {
2152 inner,
2153 offsets,
2154 list_name,
2155 list_nullable,
2156 } => {
2157 let vals = inner.finish();
2158 std::mem::replace(offsets, ListOffsets::new(false)).finish(
2159 vals,
2160 list_name,
2161 *list_nullable,
2162 )
2163 }
2164 Self::Map {
2165 key_decoder,
2166 value_decoder,
2167 offsets,
2168 map_value_name,
2169 map_value_nullable,
2170 } => {
2171 let (_, key_array) = key_decoder.finish(false);
2172 let (_, value_array) = value_decoder.finish(*map_value_nullable);
2173 let key_field = Arc::new(Field::new("key", key_array.data_type().clone(), false));
2174 let value_field = Arc::new(Field::new(
2175 &**map_value_name,
2176 value_array.data_type().clone(),
2177 *map_value_nullable,
2178 ));
2179 let entries_struct_type = DataType::Struct(
2180 vec![key_field.as_ref().clone(), value_field.as_ref().clone()].into(),
2181 );
2182 let entry_struct =
2183 StructArray::from(vec![(key_field, key_array), (value_field, value_array)]);
2184 let map_dt = DataType::Map(
2185 Arc::new(Field::new("entries", entries_struct_type, false)),
2186 false,
2187 );
2188 let len = offsets.len() - 1;
2189 let offsets_buf = Buffer::from_vec(std::mem::take(offsets));
2190 let map_data = ArrayData::builder(map_dt)
2191 .len(len)
2192 .add_buffer(offsets_buf)
2193 .add_child_data(entry_struct.into_data())
2194 .build()
2195 .unwrap();
2196 Arc::new(MapArray::from(map_data))
2197 }
2198 };
2199 let field = Field::new("", array.data_type().clone(), nullable);
2200 (field, array)
2201 }
2202}
2203
2204fn decode_repeated_varint<T: ArrowPrimitiveType>(
2209 wire_type: u8,
2210 buf: &[u8],
2211 builder: &mut PrimitiveBuilder<T>,
2212 convert: fn(u64) -> T::Native,
2213) -> Result<usize, prost::DecodeError> {
2214 if wire_type == 2 {
2215 let (data, total) = read_length_delimited(buf)?;
2217 let mut p = 0;
2218 while p < data.len() {
2219 let (v, n) = decode_varint(&data[p..])?;
2220 builder.append_value(convert(v));
2221 p += n;
2222 }
2223 Ok(total)
2224 } else if wire_type == 0 {
2225 let (v, n) = decode_varint(buf)?;
2226 builder.append_value(convert(v));
2227 Ok(n)
2228 } else {
2229 skip_field(wire_type, buf)
2230 }
2231}
2232
2233fn decode_repeated_fixed<T: ArrowPrimitiveType, const WIDTH: usize>(
2234 wire_type: u8,
2235 expected_wt: u8,
2236 buf: &[u8],
2237 builder: &mut PrimitiveBuilder<T>,
2238 convert: fn([u8; WIDTH]) -> T::Native,
2239) -> Result<usize, prost::DecodeError> {
2240 if wire_type == 2 {
2241 let (data, total) = read_length_delimited(buf)?;
2242 let mut p = 0;
2243 while p + WIDTH <= data.len() {
2244 let mut bytes = [0u8; WIDTH];
2245 bytes.copy_from_slice(&data[p..p + WIDTH]);
2246 builder.append_value(convert(bytes));
2247 p += WIDTH;
2248 }
2249 Ok(total)
2250 } else if wire_type == expected_wt {
2251 if buf.len() < WIDTH {
2252 return Err(decode_error("unexpected EOF"));
2253 }
2254 let mut bytes = [0u8; WIDTH];
2255 bytes.copy_from_slice(&buf[..WIDTH]);
2256 builder.append_value(convert(bytes));
2257 Ok(WIDTH)
2258 } else {
2259 skip_field(wire_type, buf)
2260 }
2261}
2262
2263fn decode_repeated_wrapper_varint<T: ArrowPrimitiveType>(
2264 wire_type: u8,
2265 buf: &[u8],
2266 builder: &mut PrimitiveBuilder<T>,
2267 convert: fn(u64) -> T::Native,
2268) -> Result<usize, prost::DecodeError> {
2269 if wire_type != 2 {
2270 return skip_field(wire_type, buf);
2271 }
2272 let (data, total) = read_length_delimited(buf)?;
2273 let (v, _) = decode_wrapper_varint(data)?;
2274 builder.append_value(convert(v));
2275 Ok(total)
2276}
2277
2278fn decode_repeated_wrapper_fixed32<T: ArrowPrimitiveType>(
2279 wire_type: u8,
2280 buf: &[u8],
2281 builder: &mut PrimitiveBuilder<T>,
2282 convert: fn([u8; 4]) -> T::Native,
2283) -> Result<usize, prost::DecodeError> {
2284 if wire_type != 2 {
2285 return skip_field(wire_type, buf);
2286 }
2287 let (data, total) = read_length_delimited(buf)?;
2288 let (bytes, _) = decode_wrapper_fixed32(data)?;
2289 builder.append_value(convert(bytes));
2290 Ok(total)
2291}
2292
2293fn decode_repeated_wrapper_fixed64<T: ArrowPrimitiveType>(
2294 wire_type: u8,
2295 buf: &[u8],
2296 builder: &mut PrimitiveBuilder<T>,
2297 convert: fn([u8; 8]) -> T::Native,
2298) -> Result<usize, prost::DecodeError> {
2299 if wire_type != 2 {
2300 return skip_field(wire_type, buf);
2301 }
2302 let (data, total) = read_length_delimited(buf)?;
2303 let (bytes, _) = decode_wrapper_fixed64(data)?;
2304 builder.append_value(convert(bytes));
2305 Ok(total)
2306}
2307
2308pub struct MessageDecoder {
2313 decoders: Vec<(FieldDecoder, FieldDescriptor)>,
2314 tag_map: Vec<Option<usize>>,
2315 list_nullable: bool,
2316 map_nullable: bool,
2317 num_rows: usize,
2318}
2319
2320impl MessageDecoder {
2321 pub fn new(descriptor: &MessageDescriptor, config: &PtarsConfig) -> Self {
2322 let mut decoders = Vec::new();
2323 let mut max_field_number: u32 = 0;
2324
2325 for field in descriptor.fields() {
2326 if let Some(decoder) = build_field_decoder(&field, config) {
2327 if field.number() > max_field_number {
2328 max_field_number = field.number();
2329 }
2330 decoders.push((decoder, field));
2331 }
2332 }
2333
2334 let mut tag_map = vec![
2335 None;
2336 if max_field_number == 0 {
2337 0
2338 } else {
2339 max_field_number as usize + 1
2340 }
2341 ];
2342 for (idx, (_, field)) in decoders.iter().enumerate() {
2343 let num = field.number() as usize;
2344 tag_map[num] = Some(idx);
2345 }
2346
2347 Self {
2348 decoders,
2349 tag_map,
2350 list_nullable: config.list_nullable,
2351 map_nullable: config.map_nullable,
2352 num_rows: 0,
2353 }
2354 }
2355
2356 fn decode_row(&mut self, buf: &[u8]) -> Result<(), prost::DecodeError> {
2357 let mut pos = 0;
2358 while pos < buf.len() {
2359 let (field_num, wire_type, n) = decode_tag(&buf[pos..])?;
2360 pos += n;
2361 let idx = if (field_num as usize) < self.tag_map.len() {
2362 self.tag_map[field_num as usize]
2363 } else {
2364 None
2365 };
2366 if let Some(idx) = idx {
2367 pos += self.decoders[idx].0.decode(wire_type, &buf[pos..])?;
2368 } else {
2369 pos += skip_field(wire_type, &buf[pos..])?;
2370 }
2371 }
2372 for (decoder, _) in &mut self.decoders {
2373 decoder.flush();
2374 }
2375 self.num_rows += 1;
2376 Ok(())
2377 }
2378
2379 fn decode_message_bytes(&mut self, buf: &[u8]) -> Result<(), prost::DecodeError> {
2382 let mut pos = 0;
2383 while pos < buf.len() {
2384 let (field_num, wire_type, n) = decode_tag(&buf[pos..])?;
2385 pos += n;
2386 let idx = if (field_num as usize) < self.tag_map.len() {
2387 self.tag_map[field_num as usize]
2388 } else {
2389 None
2390 };
2391 if let Some(idx) = idx {
2392 pos += self.decoders[idx].0.decode(wire_type, &buf[pos..])?;
2393 } else {
2394 pos += skip_field(wire_type, &buf[pos..])?;
2395 }
2396 }
2397 Ok(())
2398 }
2399
2400 fn flush_row(&mut self) {
2401 for (decoder, _) in &mut self.decoders {
2402 decoder.flush();
2403 }
2404 }
2405
2406 fn flush_defaults(&mut self) {
2407 for (decoder, _) in &mut self.decoders {
2408 decoder.flush(); }
2410 }
2411
2412 fn decode_null_row(&mut self) {
2413 for (decoder, _) in &mut self.decoders {
2414 decoder.flush();
2415 }
2416 self.num_rows += 1;
2417 }
2418
2419 fn row_count(&self) -> usize {
2420 self.num_rows
2421 }
2422
2423 fn build_struct_array(&mut self, validity: Option<arrow_array::BooleanArray>) -> StructArray {
2424 if self.decoders.is_empty() {
2425 let len = validity.as_ref().map_or(self.num_rows, |v| v.len());
2426 return StructArray::new_empty_fields(
2427 len,
2428 validity.map(|v| arrow::buffer::NullBuffer::new(v.values().clone())),
2429 );
2430 }
2431
2432 let (fields, columns): (Vec<_>, Vec<_>) = self
2433 .decoders
2434 .iter_mut()
2435 .map(|(decoder, field_desc)| {
2436 let nullable = if field_desc.is_list() {
2437 self.list_nullable
2438 } else if field_desc.is_map() {
2439 self.map_nullable
2440 } else {
2441 field_desc.supports_presence()
2442 };
2443 let (_, array) = decoder.finish(nullable);
2444 let field = Field::new(field_desc.name(), array.data_type().clone(), nullable);
2445 (field, array)
2446 })
2447 .unzip();
2448
2449 StructArray::new(
2450 arrow_schema::Fields::from(fields),
2451 columns,
2452 validity.map(|v| arrow::buffer::NullBuffer::new(v.values().clone())),
2453 )
2454 }
2455
2456 pub fn finish(mut self) -> RecordBatch {
2457 if self.decoders.is_empty() {
2458 let schema = Arc::new(arrow_schema::Schema::empty());
2459 return RecordBatch::try_new_with_options(
2460 schema,
2461 vec![],
2462 &arrow_array::RecordBatchOptions::new().with_row_count(Some(self.num_rows)),
2463 )
2464 .unwrap();
2465 }
2466 let struct_array = self.build_struct_array(None);
2467 RecordBatch::from(struct_array)
2468 }
2469}
2470
2471fn build_field_decoder(field: &FieldDescriptor, config: &PtarsConfig) -> Option<FieldDecoder> {
2476 if field.is_map() {
2477 return build_map_decoder(field, config);
2478 }
2479 if field.is_list() {
2480 return build_repeated_decoder(field, config);
2481 }
2482
2483 let has_presence = field.supports_presence();
2484 match field.kind() {
2485 Kind::Int32 => Some(FieldDecoder::Int32 {
2486 value: 0,
2487 has_value: false,
2488 has_presence,
2489 builder: PrimitiveBuilder::new(),
2490 }),
2491 Kind::Int64 => Some(FieldDecoder::Int64 {
2492 value: 0,
2493 has_value: false,
2494 has_presence,
2495 builder: PrimitiveBuilder::new(),
2496 }),
2497 Kind::Uint32 => Some(FieldDecoder::UInt32 {
2498 value: 0,
2499 has_value: false,
2500 has_presence,
2501 builder: PrimitiveBuilder::new(),
2502 }),
2503 Kind::Uint64 => Some(FieldDecoder::UInt64 {
2504 value: 0,
2505 has_value: false,
2506 has_presence,
2507 builder: PrimitiveBuilder::new(),
2508 }),
2509 Kind::Sint32 => Some(FieldDecoder::Sint32 {
2510 value: 0,
2511 has_value: false,
2512 has_presence,
2513 builder: PrimitiveBuilder::new(),
2514 }),
2515 Kind::Sint64 => Some(FieldDecoder::Sint64 {
2516 value: 0,
2517 has_value: false,
2518 has_presence,
2519 builder: PrimitiveBuilder::new(),
2520 }),
2521 Kind::Sfixed32 => Some(FieldDecoder::Sfixed32 {
2522 value: 0,
2523 has_value: false,
2524 has_presence,
2525 builder: PrimitiveBuilder::new(),
2526 }),
2527 Kind::Sfixed64 => Some(FieldDecoder::Sfixed64 {
2528 value: 0,
2529 has_value: false,
2530 has_presence,
2531 builder: PrimitiveBuilder::new(),
2532 }),
2533 Kind::Fixed32 => Some(FieldDecoder::Fixed32 {
2534 value: 0,
2535 has_value: false,
2536 has_presence,
2537 builder: PrimitiveBuilder::new(),
2538 }),
2539 Kind::Fixed64 => Some(FieldDecoder::Fixed64 {
2540 value: 0,
2541 has_value: false,
2542 has_presence,
2543 builder: PrimitiveBuilder::new(),
2544 }),
2545 Kind::Float => Some(FieldDecoder::Float {
2546 value: 0.0,
2547 has_value: false,
2548 has_presence,
2549 builder: PrimitiveBuilder::new(),
2550 }),
2551 Kind::Double => Some(FieldDecoder::Double {
2552 value: 0.0,
2553 has_value: false,
2554 has_presence,
2555 builder: PrimitiveBuilder::new(),
2556 }),
2557 Kind::Bool => Some(FieldDecoder::Bool {
2558 value: false,
2559 has_value: false,
2560 has_presence,
2561 builder: BooleanBuilder::new(),
2562 }),
2563 Kind::String => Some(FieldDecoder::String {
2564 value: Vec::new(),
2565 has_value: false,
2566 has_presence,
2567 builder: StringBuilderInner::new(config.use_large_string),
2568 }),
2569 Kind::Bytes => Some(FieldDecoder::Bytes {
2570 value: Vec::new(),
2571 has_value: false,
2572 has_presence,
2573 builder: BinaryBuilderInner::new(config.use_large_binary),
2574 }),
2575 Kind::Enum(enum_desc) => match config.enum_repr {
2576 EnumRepr::Int32 => Some(FieldDecoder::EnumInt32 {
2577 value: 0,
2578 has_value: false,
2579 has_presence,
2580 builder: PrimitiveBuilder::new(),
2581 }),
2582 EnumRepr::String => Some(FieldDecoder::EnumString {
2583 value: 0,
2584 has_value: false,
2585 has_presence,
2586 builder: StringBuilderInner::new(config.use_large_string),
2587 enum_descriptor: enum_desc,
2588 }),
2589 EnumRepr::Binary => Some(FieldDecoder::EnumBinary {
2590 value: 0,
2591 has_value: false,
2592 has_presence,
2593 builder: BinaryBuilderInner::new(config.use_large_binary),
2594 enum_descriptor: enum_desc,
2595 }),
2596 },
2597 Kind::Message(msg_desc) => build_message_field_decoder(msg_desc, config),
2598 }
2599}
2600
2601fn build_message_field_decoder(
2602 msg_desc: MessageDescriptor,
2603 config: &PtarsConfig,
2604) -> Option<FieldDecoder> {
2605 match msg_desc.full_name() {
2606 "google.protobuf.Timestamp" => Some(FieldDecoder::Timestamp {
2607 seconds: 0,
2608 nanos: 0,
2609 has_value: false,
2610 builder: PrimitiveBuilder::new(),
2611 unit: config.timestamp_unit.into(),
2612 tz: config.timestamp_tz.clone(),
2613 }),
2614 "google.protobuf.Duration" => Some(FieldDecoder::Duration {
2615 seconds: 0,
2616 nanos: 0,
2617 has_value: false,
2618 builder: PrimitiveBuilder::new(),
2619 unit: config.duration_unit.into(),
2620 }),
2621 "google.type.Date" => Some(FieldDecoder::Date {
2622 year: 0,
2623 month: 0,
2624 day: 0,
2625 has_value: false,
2626 builder: PrimitiveBuilder::new(),
2627 }),
2628 "google.type.TimeOfDay" => Some(FieldDecoder::TimeOfDay {
2629 hours: 0,
2630 minutes: 0,
2631 seconds_val: 0,
2632 nanos: 0,
2633 has_value: false,
2634 builder: PrimitiveBuilder::new(),
2635 unit: config.time_unit.into(),
2636 }),
2637 "google.protobuf.DoubleValue" => Some(FieldDecoder::WrapperDouble {
2638 value: 0.0,
2639 has_value: false,
2640 builder: PrimitiveBuilder::new(),
2641 }),
2642 "google.protobuf.FloatValue" => Some(FieldDecoder::WrapperFloat {
2643 value: 0.0,
2644 has_value: false,
2645 builder: PrimitiveBuilder::new(),
2646 }),
2647 "google.protobuf.Int64Value" => Some(FieldDecoder::WrapperInt64 {
2648 value: 0,
2649 has_value: false,
2650 builder: PrimitiveBuilder::new(),
2651 }),
2652 "google.protobuf.UInt64Value" => Some(FieldDecoder::WrapperUInt64 {
2653 value: 0,
2654 has_value: false,
2655 builder: PrimitiveBuilder::new(),
2656 }),
2657 "google.protobuf.Int32Value" => Some(FieldDecoder::WrapperInt32 {
2658 value: 0,
2659 has_value: false,
2660 builder: PrimitiveBuilder::new(),
2661 }),
2662 "google.protobuf.UInt32Value" => Some(FieldDecoder::WrapperUInt32 {
2663 value: 0,
2664 has_value: false,
2665 builder: PrimitiveBuilder::new(),
2666 }),
2667 "google.protobuf.BoolValue" => Some(FieldDecoder::WrapperBool {
2668 value: false,
2669 has_value: false,
2670 builder: BooleanBuilder::new(),
2671 }),
2672 "google.protobuf.StringValue" => Some(FieldDecoder::WrapperString {
2673 value: Vec::new(),
2674 has_value: false,
2675 builder: StringBuilderInner::new(config.use_large_string),
2676 }),
2677 "google.protobuf.BytesValue" => Some(FieldDecoder::WrapperBytes {
2678 value: Vec::new(),
2679 has_value: false,
2680 builder: BinaryBuilderInner::new(config.use_large_binary),
2681 }),
2682 _ => {
2683 let sub_decoder = MessageDecoder::new(&msg_desc, config);
2684 Some(FieldDecoder::Message {
2685 sub_decoder,
2686 has_value: false,
2687 is_valid: BooleanBuilder::new(),
2688 })
2689 }
2690 }
2691}
2692
2693fn build_repeated_decoder(field: &FieldDescriptor, config: &PtarsConfig) -> Option<FieldDecoder> {
2694 let ln = config.list_value_name.clone();
2695 let lnb = config.list_value_nullable;
2696 let offsets = || ListOffsets::new(config.use_large_list);
2697
2698 let inner = match field.kind() {
2699 Kind::Int32 => RepeatedInner::Int32 {
2700 values_builder: PrimitiveBuilder::new(),
2701 },
2702 Kind::Sint32 => RepeatedInner::Sint32 {
2703 values_builder: PrimitiveBuilder::new(),
2704 },
2705 Kind::Sfixed32 => RepeatedInner::Sfixed32 {
2706 values_builder: PrimitiveBuilder::new(),
2707 },
2708 Kind::Int64 => RepeatedInner::Int64 {
2709 values_builder: PrimitiveBuilder::new(),
2710 },
2711 Kind::Sint64 => RepeatedInner::Sint64 {
2712 values_builder: PrimitiveBuilder::new(),
2713 },
2714 Kind::Sfixed64 => RepeatedInner::Sfixed64 {
2715 values_builder: PrimitiveBuilder::new(),
2716 },
2717 Kind::Uint32 => RepeatedInner::UInt32 {
2718 values_builder: PrimitiveBuilder::new(),
2719 },
2720 Kind::Fixed32 => RepeatedInner::Fixed32 {
2721 values_builder: PrimitiveBuilder::new(),
2722 },
2723 Kind::Uint64 => RepeatedInner::UInt64 {
2724 values_builder: PrimitiveBuilder::new(),
2725 },
2726 Kind::Fixed64 => RepeatedInner::Fixed64 {
2727 values_builder: PrimitiveBuilder::new(),
2728 },
2729 Kind::Float => RepeatedInner::Float {
2730 values_builder: PrimitiveBuilder::new(),
2731 },
2732 Kind::Double => RepeatedInner::Double {
2733 values_builder: PrimitiveBuilder::new(),
2734 },
2735 Kind::Bool => RepeatedInner::Bool {
2736 values_builder: BooleanBuilder::new(),
2737 },
2738 Kind::String => RepeatedInner::String {
2739 values_builder: StringBuilderInner::new(config.use_large_string),
2740 },
2741 Kind::Bytes => RepeatedInner::Bytes {
2742 values_builder: BinaryBuilderInner::new(config.use_large_binary),
2743 },
2744 Kind::Enum(enum_desc) => match config.enum_repr {
2745 EnumRepr::Int32 => RepeatedInner::EnumInt32 {
2746 values_builder: PrimitiveBuilder::new(),
2747 },
2748 EnumRepr::String => RepeatedInner::EnumString {
2749 values_builder: StringBuilderInner::new(config.use_large_string),
2750 enum_descriptor: enum_desc,
2751 },
2752 EnumRepr::Binary => RepeatedInner::EnumBinary {
2753 values_builder: BinaryBuilderInner::new(config.use_large_binary),
2754 enum_descriptor: enum_desc,
2755 },
2756 },
2757 Kind::Message(msg_desc) => {
2758 return build_repeated_message_decoder(&msg_desc, config, offsets(), ln, lnb);
2759 }
2760 };
2761
2762 Some(FieldDecoder::Repeated {
2763 inner,
2764 offsets: offsets(),
2765 list_name: ln,
2766 list_nullable: lnb,
2767 })
2768}
2769
2770fn build_repeated_message_decoder(
2771 msg_desc: &MessageDescriptor,
2772 config: &PtarsConfig,
2773 offsets: ListOffsets,
2774 ln: Arc<str>,
2775 lnb: bool,
2776) -> Option<FieldDecoder> {
2777 let inner = match msg_desc.full_name() {
2778 "google.protobuf.Timestamp" => RepeatedInner::Timestamp {
2779 values_builder: PrimitiveBuilder::new(),
2780 unit: config.timestamp_unit.into(),
2781 tz: config.timestamp_tz.clone(),
2782 },
2783 "google.protobuf.Duration" => RepeatedInner::Duration {
2784 values_builder: PrimitiveBuilder::new(),
2785 unit: config.duration_unit.into(),
2786 },
2787 "google.type.Date" => RepeatedInner::Date {
2788 values_builder: PrimitiveBuilder::new(),
2789 },
2790 "google.type.TimeOfDay" => RepeatedInner::TimeOfDay {
2791 values_builder: PrimitiveBuilder::new(),
2792 unit: config.time_unit.into(),
2793 },
2794 "google.protobuf.DoubleValue" => RepeatedInner::WrapperDouble {
2795 values_builder: PrimitiveBuilder::new(),
2796 },
2797 "google.protobuf.FloatValue" => RepeatedInner::WrapperFloat {
2798 values_builder: PrimitiveBuilder::new(),
2799 },
2800 "google.protobuf.Int64Value" => RepeatedInner::WrapperInt64 {
2801 values_builder: PrimitiveBuilder::new(),
2802 },
2803 "google.protobuf.UInt64Value" => RepeatedInner::WrapperUInt64 {
2804 values_builder: PrimitiveBuilder::new(),
2805 },
2806 "google.protobuf.Int32Value" => RepeatedInner::WrapperInt32 {
2807 values_builder: PrimitiveBuilder::new(),
2808 },
2809 "google.protobuf.UInt32Value" => RepeatedInner::WrapperUInt32 {
2810 values_builder: PrimitiveBuilder::new(),
2811 },
2812 "google.protobuf.BoolValue" => RepeatedInner::WrapperBool {
2813 values_builder: BooleanBuilder::new(),
2814 },
2815 "google.protobuf.StringValue" => RepeatedInner::WrapperString {
2816 values_builder: StringBuilderInner::new(config.use_large_string),
2817 },
2818 "google.protobuf.BytesValue" => RepeatedInner::WrapperBytes {
2819 values_builder: BinaryBuilderInner::new(config.use_large_binary),
2820 },
2821 _ => {
2822 let sub_decoder = MessageDecoder::new(msg_desc, config);
2823 RepeatedInner::Message { sub_decoder }
2824 }
2825 };
2826
2827 Some(FieldDecoder::Repeated {
2828 inner,
2829 offsets,
2830 list_name: ln,
2831 list_nullable: lnb,
2832 })
2833}
2834
2835fn build_map_decoder(field: &FieldDescriptor, config: &PtarsConfig) -> Option<FieldDecoder> {
2836 let map_entry = match field.kind() {
2837 Kind::Message(desc) => desc,
2838 _ => return None,
2839 };
2840 let key_field = map_entry.get_field_by_name("key")?;
2841 let value_field = map_entry.get_field_by_name("value")?;
2842
2843 let key_decoder = build_singular_decoder_for_map(&key_field, config)?;
2845 let value_decoder = build_singular_decoder_for_map(&value_field, config)?;
2846
2847 Some(FieldDecoder::Map {
2848 key_decoder: Box::new(key_decoder),
2849 value_decoder: Box::new(value_decoder),
2850 offsets: vec![0],
2851 map_value_name: config.map_value_name.clone(),
2852 map_value_nullable: config.map_value_nullable,
2853 })
2854}
2855
2856fn build_singular_decoder_for_map(
2858 field: &FieldDescriptor,
2859 config: &PtarsConfig,
2860) -> Option<FieldDecoder> {
2861 match field.kind() {
2863 Kind::Int32 => Some(FieldDecoder::Int32 {
2864 value: 0,
2865 has_value: false,
2866 has_presence: false,
2867 builder: PrimitiveBuilder::new(),
2868 }),
2869 Kind::Int64 => Some(FieldDecoder::Int64 {
2870 value: 0,
2871 has_value: false,
2872 has_presence: false,
2873 builder: PrimitiveBuilder::new(),
2874 }),
2875 Kind::Uint32 => Some(FieldDecoder::UInt32 {
2876 value: 0,
2877 has_value: false,
2878 has_presence: false,
2879 builder: PrimitiveBuilder::new(),
2880 }),
2881 Kind::Uint64 => Some(FieldDecoder::UInt64 {
2882 value: 0,
2883 has_value: false,
2884 has_presence: false,
2885 builder: PrimitiveBuilder::new(),
2886 }),
2887 Kind::Sint32 => Some(FieldDecoder::Sint32 {
2888 value: 0,
2889 has_value: false,
2890 has_presence: false,
2891 builder: PrimitiveBuilder::new(),
2892 }),
2893 Kind::Sint64 => Some(FieldDecoder::Sint64 {
2894 value: 0,
2895 has_value: false,
2896 has_presence: false,
2897 builder: PrimitiveBuilder::new(),
2898 }),
2899 Kind::Sfixed32 => Some(FieldDecoder::Sfixed32 {
2900 value: 0,
2901 has_value: false,
2902 has_presence: false,
2903 builder: PrimitiveBuilder::new(),
2904 }),
2905 Kind::Sfixed64 => Some(FieldDecoder::Sfixed64 {
2906 value: 0,
2907 has_value: false,
2908 has_presence: false,
2909 builder: PrimitiveBuilder::new(),
2910 }),
2911 Kind::Fixed32 => Some(FieldDecoder::Fixed32 {
2912 value: 0,
2913 has_value: false,
2914 has_presence: false,
2915 builder: PrimitiveBuilder::new(),
2916 }),
2917 Kind::Fixed64 => Some(FieldDecoder::Fixed64 {
2918 value: 0,
2919 has_value: false,
2920 has_presence: false,
2921 builder: PrimitiveBuilder::new(),
2922 }),
2923 Kind::Float => Some(FieldDecoder::Float {
2924 value: 0.0,
2925 has_value: false,
2926 has_presence: false,
2927 builder: PrimitiveBuilder::new(),
2928 }),
2929 Kind::Double => Some(FieldDecoder::Double {
2930 value: 0.0,
2931 has_value: false,
2932 has_presence: false,
2933 builder: PrimitiveBuilder::new(),
2934 }),
2935 Kind::Bool => Some(FieldDecoder::Bool {
2936 value: false,
2937 has_value: false,
2938 has_presence: false,
2939 builder: BooleanBuilder::new(),
2940 }),
2941 Kind::String => Some(FieldDecoder::String {
2942 value: Vec::new(),
2943 has_value: false,
2944 has_presence: false,
2945 builder: StringBuilderInner::new(config.use_large_string),
2946 }),
2947 Kind::Bytes => Some(FieldDecoder::Bytes {
2948 value: Vec::new(),
2949 has_value: false,
2950 has_presence: false,
2951 builder: BinaryBuilderInner::new(config.use_large_binary),
2952 }),
2953 Kind::Enum(enum_desc) => match config.enum_repr {
2954 EnumRepr::Int32 => Some(FieldDecoder::EnumInt32 {
2955 value: 0,
2956 has_value: false,
2957 has_presence: false,
2958 builder: PrimitiveBuilder::new(),
2959 }),
2960 EnumRepr::String => Some(FieldDecoder::EnumString {
2961 value: 0,
2962 has_value: false,
2963 has_presence: false,
2964 builder: StringBuilderInner::new(config.use_large_string),
2965 enum_descriptor: enum_desc,
2966 }),
2967 EnumRepr::Binary => Some(FieldDecoder::EnumBinary {
2968 value: 0,
2969 has_value: false,
2970 has_presence: false,
2971 builder: BinaryBuilderInner::new(config.use_large_binary),
2972 enum_descriptor: enum_desc,
2973 }),
2974 },
2975 Kind::Message(msg_desc) => build_message_field_decoder(msg_desc, config),
2976 }
2977}
2978
2979pub fn binary_array_to_record_batch_direct(
2988 array: &BinaryArray,
2989 descriptor: &MessageDescriptor,
2990 config: &PtarsConfig,
2991) -> Result<RecordBatch, prost::DecodeError> {
2992 let mut decoder = MessageDecoder::new(descriptor, config);
2993 let policy = config.confluent_wire_policy;
2994 for i in 0..array.len() {
2995 if array.is_null(i) {
2996 decoder.decode_null_row();
2997 } else {
2998 let bytes = strip_confluent_prefix(array.value(i), policy)?;
2999 decoder.decode_row(bytes)?;
3000 }
3001 }
3002 Ok(decoder.finish())
3003}
3004
3005pub fn messages_to_record_batch(
3010 messages: &[prost_reflect::DynamicMessage],
3011 message_descriptor: &MessageDescriptor,
3012) -> RecordBatch {
3013 messages_to_record_batch_with_config(messages, message_descriptor, &PtarsConfig::default())
3014}
3015
3016pub fn messages_to_record_batch_with_config(
3021 messages: &[prost_reflect::DynamicMessage],
3022 message_descriptor: &MessageDescriptor,
3023 config: &PtarsConfig,
3024) -> RecordBatch {
3025 use arrow_array::builder::BinaryBuilder;
3026 use prost::Message;
3027
3028 let mut bin_builder = BinaryBuilder::new();
3029 for msg in messages {
3030 bin_builder.append_value(msg.encode_to_vec());
3031 }
3032 let binary_array = bin_builder.finish();
3033 binary_array_to_record_batch_direct(&binary_array, message_descriptor, config)
3034 .expect("failed to decode messages")
3035}
3036
3037pub fn binary_array_to_messages(
3042 array: &BinaryArray,
3043 message_descriptor: &MessageDescriptor,
3044) -> Result<Vec<prost_reflect::DynamicMessage>, prost::DecodeError> {
3045 let mut messages = Vec::with_capacity(array.len());
3046 for i in 0..array.len() {
3047 let message = if array.is_null(i) {
3048 prost_reflect::DynamicMessage::new(message_descriptor.clone())
3049 } else {
3050 prost_reflect::DynamicMessage::decode(message_descriptor.clone(), array.value(i))?
3051 };
3052 messages.push(message);
3053 }
3054 Ok(messages)
3055}
3056
3057#[cfg(test)]
3058mod tests {
3059 use super::*;
3060
3061 #[test]
3062 fn test_strip_confluent_prefix_raw() {
3063 let buf = b"\x00\x00\x00\x00\x01\x08\x96\x01";
3064 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Raw).unwrap();
3065 assert_eq!(result, buf);
3066 }
3067
3068 #[test]
3069 fn test_strip_confluent_prefix_standard() {
3070 let buf = b"\x00\x00\x00\x00\x01\x08\x96\x01";
3072 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Standard).unwrap();
3073 assert_eq!(result, b"\x08\x96\x01");
3074 }
3075
3076 #[test]
3077 fn test_strip_confluent_prefix_standard_too_short() {
3078 let buf = b"\x00\x01\x02";
3079 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Standard);
3080 assert!(result.is_err());
3081 }
3082
3083 #[test]
3084 fn test_strip_confluent_prefix_protobuf_zero_indexes() {
3085 let buf = b"\x00\x00\x00\x00\x01\x00\x08\x96\x01";
3087 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Protobuf).unwrap();
3088 assert_eq!(result, b"\x08\x96\x01");
3089 }
3090
3091 #[test]
3092 fn test_strip_confluent_prefix_protobuf_one_index() {
3093 let buf = b"\x00\x00\x00\x00\x01\x01\x00\x08\x96\x01";
3095 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Protobuf).unwrap();
3096 assert_eq!(result, b"\x08\x96\x01");
3097 }
3098
3099 #[test]
3100 fn test_strip_confluent_prefix_protobuf_two_indexes() {
3101 let buf = b"\x00\x00\x00\x00\x01\x02\x04\x02\x08\x96\x01";
3103 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Protobuf).unwrap();
3104 assert_eq!(result, b"\x08\x96\x01");
3105 }
3106
3107 #[test]
3108 fn test_strip_confluent_prefix_protobuf_too_short() {
3109 let buf = b"\x00\x01";
3110 let result = strip_confluent_prefix(buf, ConfluentWirePolicy::Protobuf);
3111 assert!(result.is_err());
3112 }
3113}