zeroshade commented on code in PR #2106: URL: https://github.com/apache/iceberg-go/pull/2106#discussion_r4198576745
########## table/readtasks_residual_binding_internal_test.go: ########## @@ -0,0 +1,195 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package table + +import ( + "context" + "testing" + + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/stretchr/testify/require" +) + +func residualBindingTestScan(t *testing.T) (*Scan, *iceberg.Schema) { + t.Helper() + + schema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64}, + ) + metadata, err := NewMetadata( + schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, "mem://mixed-residuals", nil, + ) + require.NoError(t, err) + memFS := iceio.NewMemFS() + tbl := New( + Identifier{"db", "tbl"}, metadata, "metadata.json", + func(context.Context) (iceio.IO, error) { return memFS, nil }, nil, + ) + + return tbl.Scan(), schema +} + +func TestBindReadTasksResidualsCopyOnWrite(t *testing.T) { + _, schema := residualBindingTestScan(t) + unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + bound, err := iceberg.BindExpr(schema, unbound, true) + require.NoError(t, err) + + tests := []struct { + name string + residuals []iceberg.BooleanExpression + wantAlias bool + }{ + {name: "all bound and nil", residuals: []iceberg.BooleanExpression{bound, nil, bound}, wantAlias: true}, + {name: "mixed", residuals: []iceberg.BooleanExpression{bound, nil, unbound, bound}}, + {name: "first task unbound", residuals: []iceberg.BooleanExpression{unbound, bound}}, + {name: "all unbound", residuals: []iceberg.BooleanExpression{unbound, unbound}}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + tasks := make([]FileScanTask, len(tt.residuals)) + for i, residual := range tt.residuals { + tasks[i].Residual = residual + } + + got, err := bindReadTasksResiduals(schema, tasks, true) + require.NoError(t, err) + require.Len(t, got, len(tasks)) + if len(tasks) > 0 { + require.Equal(t, tt.wantAlias, &got[0] == &tasks[0]) + } + + for i, original := range tt.residuals { + if original == nil { + require.Nil(t, got[i].Residual) + + continue + } + + // The input plan is never rewritten, even when the output needs binding. + require.Same(t, original, tasks[i].Residual) + state, visitErr := iceberg.VisitExpr(got[i].Residual, filterBindingVisitor{}) + require.NoError(t, visitErr) + require.True(t, state.hasBound) + require.False(t, state.hasUnbound) + if stateBefore, stateErr := iceberg.VisitExpr(original, filterBindingVisitor{}); stateErr == nil && !stateBefore.hasUnbound { + require.Same(t, original, got[i].Residual) + } + } + }) + } +} + +func TestReadTasksResidualPlanIsReusable(t *testing.T) { + scan, schema := residualBindingTestScan(t) + unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + bound, err := iceberg.BindExpr(schema, unbound, true) + require.NoError(t, err) + + wrongSchema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 2, Name: "other", Type: iceberg.PrimitiveTypes.Int64}, + ) + wrongBound, err := iceberg.BindExpr( + wrongSchema, iceberg.EqualTo(iceberg.Reference("other"), int64(1)), true, + ) + require.NoError(t, err) + + tests := []struct { + name string + tasks []FileScanTask + wantErr bool + }{ + { + name: "mixed plan", + tasks: []FileScanTask{{Residual: bound}, {}, {Residual: unbound}, {Residual: bound}}, + }, + { + name: "invalid bound residual", + tasks: []FileScanTask{{Residual: bound}, {Residual: unbound}, {Residual: wrongBound}}, + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + originals := make([]iceberg.BooleanExpression, len(tt.tasks)) + for i := range tt.tasks { + originals[i] = tt.tasks[i].Residual + } + + // Run twice to prove the caller-owned plan remains reusable. + for range 2 { + _, _, err := scan.ReadTasks(t.Context(), tt.tasks) + if tt.wantErr { + require.ErrorIs(t, err, iceberg.ErrInvalidArgument) + require.ErrorContains(t, err, "field ID 2") + } else { + require.NoError(t, err) + } + for i, original := range originals { + if original == nil { + require.Nil(t, tt.tasks[i].Residual) + } else { + require.Same(t, original, tt.tasks[i].Residual) + } + } + } + }) + } +} + +func TestReadTasksAlreadyBoundTasksRemainReadOnly(t *testing.T) { + scan, schema := residualBindingTestScan(t) + unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + bound, err := iceberg.BindExpr(schema, unbound, true) + require.NoError(t, err) + + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, + iceberg.EntryContentData, + "mem://mixed-residuals/missing.parquet", + iceberg.ParquetFile, + nil, + nil, + nil, + 1, + 128, + ) + require.NoError(t, err) + file := builder.Build() + tasks := []FileScanTask{{File: file, Residual: bound}} + original := tasks[0] + + _, records, err := scan.ReadTasks(t.Context(), tasks) + require.NoError(t, err) + + sawReadError := false + for _, readErr := range records { + require.Error(t, readErr) + sawReadError = true + + break + } + require.True(t, sawReadError) + require.Same(t, original.File, tasks[0].File) + require.Same(t, original.Residual, tasks[0].Residual) + require.Equal(t, original.Start, tasks[0].Start) + require.Equal(t, original.Length, tasks[0].Length) Review Comment: This is the test meant to pin "downstream only reads", but it only compares four fields after the first missing-file error, so it misses two plausible regressions: - A write that stores an equal value. For example, caching the per-task filter back via `tasks[i].Residual = rowFilter`: `rowFilterForTask` returns the same pointer for an already-bound residual, so `require.Same` still passes. - A write to a field it doesn't check (`DeleteFiles`, `DeletionVectorFiles`, `FirstRowID`, ...). CI runs `make test-race`, so a stronger guard is cheap: 1. Append a couple of real data files and put an already-bound residual on every task. 2. Drain two `ReadTasks` iterators concurrently over the same slice, one scan per goroutine, with `WithMaxConcurrency(4)`. 3. Finish with `require.Equal(t, snapshot, tasks)`. Any write into the aliased slice from `GetRecords` setup or the feeder goroutine then fails under `-race`, and the whole-struct comparison covers the remaining fields. I ran that shape as a throwaway against this head and it's clean. ########## table/scanner.go: ########## @@ -2161,6 +2193,10 @@ func (scan *Scan) ToArrowRecords(ctx context.Context) (*arrow.Schema, iter.Seq2[ // reached; if no such task is processed, the file is not read and its error is not // returned. The returned iterator is single-use. // +// The caller must not modify tasks or any task element until the returned +// iterator is exhausted or abandoned. When every residual is already bound, +// ReadTasks may retain the caller's backing array instead of cloning it. Review Comment: Nit: tasks with a nil residual take the alias path too (that's the common unfiltered `ToArrowRecords` case), so "every residual is already bound" undersells when the caller's array is retained. ```suggestion // The caller must not modify tasks or any task element until the returned // iterator is exhausted or abandoned. When no residual needs binding (each is // nil or already bound), ReadTasks may retain the caller's backing array // instead of cloning it. ``` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
