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 4278e866 fix(arrow/ipc): release reader state on schema errors (#1112)
4278e866 is described below

commit 4278e8664f006c5481a72a29affd3969d84ae4b9
Author: Minh Vu <[email protected]>
AuthorDate: Fri Aug 7 20:43:08 2026 +0200

    fix(arrow/ipc): release reader state on schema errors (#1112)
    
    ### Rationale for this change
    
    The IPC reader constructors retain reader state before schema
    initialization. When schema loading fails, they can return without
    releasing the message reader.
    
    ### What changes are included in this PR?
    
    Release the reader state on schema errors and recovered panics, and add
    a message reader that verifies the release call.
    
    ### Are these changes tested?
    
    - `go test ./arrow/ipc`
    
    ### Are there any user-facing changes?
    
    No API changes. This corrects the reported behavior while preserving the
    existing ownership and compatibility contracts.
---
 arrow/ipc/ipc_test.go | 33 +++++++++++++++++++++++++++++++++
 arrow/ipc/reader.go   | 18 ++++++++++++++----
 2 files changed, 47 insertions(+), 4 deletions(-)

diff --git a/arrow/ipc/ipc_test.go b/arrow/ipc/ipc_test.go
index 93bd987a..b7e006ba 100644
--- a/arrow/ipc/ipc_test.go
+++ b/arrow/ipc/ipc_test.go
@@ -132,6 +132,39 @@ func (r *testMessageReader) Message() (*ipc.Message, 
error) {
 func (r *testMessageReader) Release() {}
 func (r *testMessageReader) Retain()  {}
 
+type releaseCountingMessageReader struct {
+       releases int
+}
+
+func (r *releaseCountingMessageReader) Message() (*ipc.Message, error) {
+       return nil, errors.New("schema read failed")
+}
+
+func (r *releaseCountingMessageReader) Release() { r.releases++ }
+func (r *releaseCountingMessageReader) Retain()  {}
+
+func TestNewReaderFromMessageReaderReleasesOnSchemaError(t *testing.T) {
+       msgReader := &releaseCountingMessageReader{}
+       _, err := ipc.NewReaderFromMessageReader(msgReader)
+       require.Error(t, err)
+       assert.Equal(t, 1, msgReader.releases)
+}
+
+type panicMessageReader struct {
+       releases int
+}
+
+func (r *panicMessageReader) Message() (*ipc.Message, error) { panic("schema 
read panicked") }
+func (r *panicMessageReader) Release()                       { r.releases++ }
+func (r *panicMessageReader) Retain()                        {}
+
+func TestNewReaderFromMessageReaderReleasesOnSchemaPanic(t *testing.T) {
+       msgReader := &panicMessageReader{}
+       _, err := ipc.NewReaderFromMessageReader(msgReader)
+       require.Error(t, err)
+       assert.Equal(t, 1, msgReader.releases)
+}
+
 // Ensure that if the MessageReader errors, we get the error from Read
 func TestArrow14769(t *testing.T) {
        reader, err := ipc.NewReaderFromMessageReader(&testMessageReader{})
diff --git a/arrow/ipc/reader.go b/arrow/ipc/reader.go
index dc035564..e15994e0 100644
--- a/arrow/ipc/reader.go
+++ b/arrow/ipc/reader.go
@@ -60,6 +60,10 @@ type Reader struct {
 func NewReaderFromMessageReader(r MessageReader, opts ...Option) (reader 
*Reader, err error) {
        defer func() {
                if pErr := recover(); pErr != nil {
+                       if reader != nil {
+                               reader.Release()
+                               reader = nil
+                       }
                        err = utils.FormatRecoveredError("arrow/ipc: unknown 
error while reading", pErr)
                }
        }()
@@ -68,7 +72,7 @@ func NewReaderFromMessageReader(r MessageReader, opts 
...Option) (reader *Reader
                opt(cfg)
        }
 
-       rr := &Reader{
+       reader = &Reader{
                r:        r,
                refCount: atomic.Int64{},
                // types:    make(dictTypeMap),
@@ -77,21 +81,26 @@ func NewReaderFromMessageReader(r MessageReader, opts 
...Option) (reader *Reader
                ensureNativeEndian: cfg.ensureNativeEndian,
                expectedSchema:     cfg.schema,
        }
-       rr.refCount.Add(1)
+       reader.refCount.Add(1)
 
        if !cfg.noAutoSchema {
-               if err := rr.readSchema(cfg.schema); err != nil {
+               if err := reader.readSchema(cfg.schema); err != nil {
+                       reader.Release()
                        return nil, err
                }
        }
 
-       return rr, nil
+       return reader, nil
 }
 
 // NewReader returns a reader that reads records from an input stream.
 func NewReader(r io.Reader, opts ...Option) (rr *Reader, err error) {
        defer func() {
                if pErr := recover(); pErr != nil {
+                       if rr != nil {
+                               rr.Release()
+                               rr = nil
+                       }
                        err = utils.FormatRecoveredError("arrow/ipc: unknown 
error while reading", pErr)
                }
        }()
@@ -111,6 +120,7 @@ func NewReader(r io.Reader, opts ...Option) (rr *Reader, 
err error) {
 
        if !cfg.noAutoSchema {
                if err := rr.readSchema(cfg.schema); err != nil {
+                       rr.Release()
                        return nil, err
                }
        }

Reply via email to