Skip to main content

ptars_core/
proto_to_arrow.rs

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// ---------------------------------------------------------------------------
22// Wire format decoding primitives
23// ---------------------------------------------------------------------------
24
25#[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
86/// Read a length-delimited field, returning (data_slice, bytes_consumed).
87fn 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
97/// Strip the Confluent Schema Registry wire format prefix from a message.
98fn 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            // Read varint-encoded count of message indexes
120            let (count, mut offset) = decode_varint(remaining)?;
121            // Skip `count` varints (the message indexes themselves)
122            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
169// ---------------------------------------------------------------------------
170// String/Binary builder inner enums
171// ---------------------------------------------------------------------------
172
173enum 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
257// ---------------------------------------------------------------------------
258// ListOffsets
259// ---------------------------------------------------------------------------
260
261enum 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
307// ---------------------------------------------------------------------------
308// RepeatedInner enum — repeated field value storage
309// ---------------------------------------------------------------------------
310
311enum 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
800// ---------------------------------------------------------------------------
801// FieldDecoder enum — all types
802// ---------------------------------------------------------------------------
803
804enum FieldDecoder {
805    // --- Singular scalars (buffered) ---
806    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    // --- Well-known types (singular, buffered) ---
918    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    // --- Wrapper types (singular) ---
951    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    // --- Nested message ---
998    Message {
999        sub_decoder: MessageDecoder,
1000        has_value: bool,
1001        is_valid: BooleanBuilder,
1002    },
1003
1004    // --- Repeated fields (all collapsed into one variant) ---
1005    Repeated {
1006        inner: RepeatedInner,
1007        offsets: ListOffsets,
1008        list_name: Arc<str>,
1009        list_nullable: bool,
1010    },
1011
1012    // --- Map fields ---
1013    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
1022// ---------------------------------------------------------------------------
1023// Helpers for decoding wire values inline
1024// ---------------------------------------------------------------------------
1025
1026/// Decode fields of a well-known submessage with up to 4 int fields.
1027/// Returns (field1, field2, field3, field4) initialized to 0, scanning for field numbers 1..=max_field.
1028fn 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
1045/// Decode a wrapper submessage: field 1 with the given wire type.
1046/// For varint wrapper types.
1047fn 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
1149// ---------------------------------------------------------------------------
1150// Macro for flush/finish boilerplate
1151// ---------------------------------------------------------------------------
1152
1153macro_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
1237// ---------------------------------------------------------------------------
1238// FieldDecoder: decode + flush + finish
1239// ---------------------------------------------------------------------------
1240
1241impl 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            // Well-known types: decode submessage
1442            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            // Wrapper types
1513            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            // Nested message
1628            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                // For singular messages, if seen multiple times, the spec says merge.
1638                // For simplicity (and matching DynamicMessage behavior), we decode fresh each time.
1639                // Reset sub_decoder if already decoded this row.
1640                if !*has_value {
1641                    *has_value = true;
1642                }
1643                sub_decoder.decode_message_bytes(data)?;
1644                Ok(total)
1645            }
1646            // Repeated fields — delegate to inner
1647            Self::Repeated { inner, .. } => inner.decode(wire_type, buf),
1648            // Map: each occurrence is a length-delimited entry submessage
1649            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                // Parse entry submessage: field 1 = key, field 2 = value
1659                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                // Flush key and value (they're buffered singular decoders)
1672                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            // Well-known types
1854            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            // Wrapper types: present → value, absent → null
1950            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            // Nested message
2068            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            // Repeated fields: push offset
2083            Self::Repeated { inner, offsets, .. } => {
2084                offsets.push(inner.len());
2085            }
2086            // Map: push offset based on key builder length
2087            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        // This is called by MessageDecoder::finish, which provides the field name separately
2110        // We return a dummy field name here; the caller replaces it.
2111        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
2204// ---------------------------------------------------------------------------
2205// Helpers for repeated varint/fixed decoding
2206// ---------------------------------------------------------------------------
2207
2208fn 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        // packed
2216        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
2308// ---------------------------------------------------------------------------
2309// MessageDecoder
2310// ---------------------------------------------------------------------------
2311
2312pub 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    /// Decode a submessage's bytes without flushing — used for singular message fields
2380    /// where the parent will call flush.
2381    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(); // all absent → defaults
2409        }
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
2471// ---------------------------------------------------------------------------
2472// build_field_decoder
2473// ---------------------------------------------------------------------------
2474
2475fn 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    // Build singular decoders for key and value (they buffer per-entry, not per-row)
2844    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
2856/// Build a singular decoder for use inside a map entry (no presence tracking).
2857fn build_singular_decoder_for_map(
2858    field: &FieldDescriptor,
2859    config: &PtarsConfig,
2860) -> Option<FieldDecoder> {
2861    // Map keys/values are never "optional" in protobuf sense — they use proto3 defaults
2862    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
2979// ---------------------------------------------------------------------------
2980// Public API
2981// ---------------------------------------------------------------------------
2982
2983/// Decode a BinaryArray of serialized protobuf messages directly into a RecordBatch.
2984///
2985/// Parses protobuf wire format directly into Arrow builders — no intermediate
2986/// message objects are created.
2987pub 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
3005/// Convert DynamicMessage instances to a RecordBatch using the default configuration.
3006///
3007/// Each message is serialized to protobuf wire format, then decoded directly
3008/// into Arrow arrays.
3009pub 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
3016/// Convert DynamicMessage instances to a RecordBatch using the specified configuration.
3017///
3018/// Each message is serialized to protobuf wire format, then decoded directly
3019/// into Arrow arrays.
3020pub 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
3037/// Decode a BinaryArray into a vector of DynamicMessage.
3038///
3039/// Each element in the binary array is expected to be a serialized protobuf message.
3040/// Null values in the array result in default (empty) messages.
3041pub 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        // magic byte + 4-byte schema ID + payload
3071        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        // magic byte + 4-byte schema ID + varint 0 (count=0) + payload
3086        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        // magic byte + 4-byte schema ID + varint 1 (count=1) + varint 0 (index) + payload
3094        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        // magic byte + 4-byte schema ID + varint 2 (count) + varint 4 + varint 2 + payload
3102        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}