This is an automated email from the ASF dual-hosted git repository.

zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-go.git


The following commit(s) were added to refs/heads/main by this push:
     new ab2efb08 fix(arrow/util): return protobuf record conversion errors 
(#1088)
ab2efb08 is described below

commit ab2efb08ff9e892d7751571da457672555306aff
Author: Minh Vu <[email protected]>
AuthorDate: Mon Aug 10 18:45:23 2026 +0200

    fix(arrow/util): return protobuf record conversion errors (#1088)
    
    ### Rationale for this change
    
    Record discards errors from AppendValueOrNull. An unknown protobuf enum
    value can then panic or leave builders with mismatched lengths instead
    of reporting the conversion failure.
    
    ### What changes are included in this PR?
    
    Add RecordWithError with field context, reject undefined enum numbers
    safely, and clean up temporary builders and dictionary values on both
    success and failure. Keep Record with its existing panic-on-error
    behavior for compatibility.
    
    ### Are these changes tested?
    
    - `go test ./arrow/util`
    - Added checked-allocator coverage for undefined enum values and
    conversion cleanup.
    
    ### Are there any user-facing changes?
    
    Yes. RecordWithError provides an error-returning path while the existing
    panic-on-error Record API remains available.
---
 arrow/util/protobuf_reflect.go      | 32 +++++++++++++++++++++++------
 arrow/util/protobuf_reflect_test.go | 40 +++++++++++++++++++++++++++++++++++++
 2 files changed, 66 insertions(+), 6 deletions(-)

diff --git a/arrow/util/protobuf_reflect.go b/arrow/util/protobuf_reflect.go
index ca17c838..e6e01109 100644
--- a/arrow/util/protobuf_reflect.go
+++ b/arrow/util/protobuf_reflect.go
@@ -661,23 +661,35 @@ func (msg ProtobufMessageReflection) Schema() 
*arrow.Schema {
 
 // Record returns an arrow.RecordBatch for a protobuf message
 func (msg ProtobufMessageReflection) Record(mem memory.Allocator) 
arrow.RecordBatch {
+       record, err := msg.RecordWithError(mem)
+       if err != nil {
+               panic(err)
+       }
+       return record
+}
+
+// RecordWithError returns an arrow.RecordBatch for a protobuf message and
+// reports conversion failures.
+func (msg ProtobufMessageReflection) RecordWithError(mem memory.Allocator) 
(arrow.RecordBatch, error) {
        if mem == nil {
                mem = memory.NewGoAllocator()
        }
 
        schema := msg.Schema()
        if schema.NumFields() == 0 {
-               return array.NewRecordBatch(schema, nil, 1)
+               return array.NewRecordBatch(schema, nil, 1), nil
        }
 
        recordBuilder := array.NewRecordBuilder(mem, schema)
        defer recordBuilder.Release()
 
        for i, f := range msg.fields {
-               f.AppendValueOrNull(recordBuilder.Field(i), mem)
+               if err := f.AppendValueOrNull(recordBuilder.Field(i), mem); err 
!= nil {
+                       return nil, fmt.Errorf("failed to append protobuf field 
%q: %w", f.name(), err)
+               }
        }
 
-       return recordBuilder.NewRecordBatch()
+       return recordBuilder.NewRecordBatch(), nil
 }
 
 // NewProtobufMessageReflection initialises a ProtobufMessageReflection
@@ -776,7 +788,11 @@ func (f ProtobufMessageFieldReflection) 
AppendValueOrNull(b array.Builder, mem m
        switch b.Type().ID() {
        case arrow.STRING:
                if f.isEnum() {
-                       
b.(*array.StringBuilder).Append(string(fd.Enum().Values().ByNumber(pv.Enum()).Name()))
+                       enumValue := fd.Enum().Values().ByNumber(pv.Enum())
+                       if enumValue == nil {
+                               return fmt.Errorf("enum value %d is not defined 
for %s", pv.Enum(), fd.Enum().FullName())
+                       }
+                       
b.(*array.StringBuilder).Append(string(enumValue.Name()))
                } else {
                        b.(*array.StringBuilder).Append(pv.String())
                }
@@ -816,6 +832,11 @@ func (f ProtobufMessageFieldReflection) 
AppendValueOrNull(b array.Builder, mem m
                        return err
                }
        case arrow.DICTIONARY:
+               enumNum := int(f.reflectValue().Int())
+               enumValue := 
fd.Enum().Values().ByNumber(protoreflect.EnumNumber(enumNum))
+               if enumValue == nil {
+                       return fmt.Errorf("enum value %d is not defined for 
%s", enumNum, fd.Enum().FullName())
+               }
                pdr := f.asDictionary()
                db := b.(*array.BinaryDictionaryBuilder)
                dictValues := pdr.getDictValues(mem).(*array.String)
@@ -824,8 +845,7 @@ func (f ProtobufMessageFieldReflection) AppendValueOrNull(b 
array.Builder, mem m
                if err != nil {
                        return err
                }
-               enumNum := int(f.reflectValue().Int())
-               enumVal := 
fd.Enum().Values().ByNumber(protoreflect.EnumNumber(enumNum)).Name()
+               enumVal := enumValue.Name()
                err = db.AppendValueFromString(string(enumVal))
                if err != nil {
                        return err
diff --git a/arrow/util/protobuf_reflect_test.go 
b/arrow/util/protobuf_reflect_test.go
index c93331cd..79cc6e28 100644
--- a/arrow/util/protobuf_reflect_test.go
+++ b/arrow/util/protobuf_reflect_test.go
@@ -390,6 +390,19 @@ func TestRecordReleasesConstructionBuffers(t *testing.T) {
        mem.AssertSize(t, 0)
 }
 
+func TestRecordWithErrorReportsFieldConversion(t *testing.T) {
+       msg := AllTheTypesNoAnyFixture().msg.(*util_message.AllTheTypesNoAny)
+       msg.Enum = util_message.AllTheTypesNoAny_ExampleEnum(999)
+       pmr := NewProtobufMessageReflection(msg)
+       mem := memory.NewCheckedAllocator(memory.NewGoAllocator())
+       rec, err := pmr.RecordWithError(mem)
+       assert.Nil(t, rec)
+       require.Error(t, err)
+       assert.ErrorContains(t, err, `failed to append protobuf field "enum"`)
+       assert.ErrorContains(t, err, "enum value 999 is not defined")
+       mem.AssertSize(t, 0)
+}
+
 func TestRecordWithAllFieldsExcluded(t *testing.T) {
        pmr := NewProtobufMessageReflection(AllTheTypesNoAnyFixture().msg,
                WithExclusionPolicy(func(*ProtobufFieldReflection) bool { 
return true }))
@@ -400,6 +413,33 @@ func TestRecordWithAllFieldsExcluded(t *testing.T) {
        assert.Zero(t, rec.NumCols())
 }
 
+func TestRecordWithErrorReportsUnknownEnumValueForEnumValueHandler(t 
*testing.T) {
+       msg := AllTheTypesNoAnyFixture().msg.(*util_message.AllTheTypesNoAny)
+       msg.Enum = util_message.AllTheTypesNoAny_ExampleEnum(999)
+       onlyEnum := func(pfr *ProtobufFieldReflection) bool { return 
!pfr.isEnum() }
+       pmr := NewProtobufMessageReflection(msg, WithExclusionPolicy(onlyEnum), 
WithEnumHandler(EnumValue))
+       mem := memory.NewCheckedAllocator(memory.NewGoAllocator())
+
+       rec, err := pmr.RecordWithError(mem)
+       assert.Nil(t, rec)
+       require.Error(t, err)
+       assert.ErrorContains(t, err, `failed to append protobuf field "enum"`)
+       assert.ErrorContains(t, err, "enum value 999 is not defined")
+       mem.AssertSize(t, 0)
+}
+
+func TestRecordWithErrorReturnsOneRowForZeroFieldSchema(t *testing.T) {
+       excludeAll := func(*ProtobufFieldReflection) bool { return true }
+       pmr := NewProtobufMessageReflection(&util_message.AllTheTypesNoAny{}, 
WithExclusionPolicy(excludeAll))
+
+       rec, err := pmr.RecordWithError(memory.DefaultAllocator)
+       require.NoError(t, err)
+       require.NotNil(t, rec)
+       defer rec.Release()
+       require.EqualValues(t, 1, rec.NumRows())
+       require.EqualValues(t, 0, rec.NumCols())
+}
+
 func TestRecordFromEmptyMessage(t *testing.T) {
        pmr := NewProtobufMessageReflection(&emptypb.Empty{})
        rec := pmr.Record(nil)

Reply via email to