Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .golangci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ linters:
- copyloopvar
- usestdlibvars
- misspell
- modernize
- nlreturn
- perfsprint
- staticcheck
Expand Down
4 changes: 2 additions & 2 deletions catalog/rest/load_table_bench_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,7 @@ func makeTableResponseWithSnapshots(snapshotCount int64) []byte {
panic(fmt.Errorf("failed to generate load table response: %w", err))
}

return []byte(fmt.Sprintf(`{
return fmt.Appendf(nil, `{
"metadata-location": "s3://warehouse/database/table/metadata/00001-5f2f8166-244c-4eae-ac36-384ecdec81fc.gz.metadata.json",
"metadata": {
"format-version": 1,
Expand Down Expand Up @@ -190,5 +190,5 @@ func makeTableResponseWithSnapshots(snapshotCount int64) []byte {
}
]
}
}`, snapshotTimestamp, snapshotID, snapshotID, snapshotsJson, snapshotsLogEntriesJson))
}`, snapshotTimestamp, snapshotID, snapshotID, snapshotsJson, snapshotsLogEntriesJson)
}
6 changes: 2 additions & 4 deletions catalog/sql/sql_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1953,12 +1953,10 @@ func (s *SqliteCatalogTestSuite) TestCreateNamespaceConcurrent() {

var wg sync.WaitGroup
for range writers {
wg.Add(1)
go func() {
defer wg.Done()
wg.Go(func() {
<-start
errs <- cat.CreateNamespace(ctx, namespace, nil)
}()
})
}
close(start)
wg.Wait()
Expand Down
2 changes: 1 addition & 1 deletion data_file_codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ func newDecodeEntry(version int) (any, *dataFile) {
return &manifestEntry{Data: df}, df
}

var dataFileAvroFieldIndexes = avroFieldIndexes(reflect.TypeOf(dataFile{}))
var dataFileAvroFieldIndexes = avroFieldIndexes(reflect.TypeFor[dataFile]())

func avroFieldIndexes(t reflect.Type) []int {
indexes := make([]int, 0, t.NumField())
Expand Down
2 changes: 1 addition & 1 deletion data_file_codec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ func TestMarshalAvroEntryDoesNotMutateAnyAvroField(t *testing.T) {
}

func TestDataFileAvroFieldIndexesCoverEveryAvroField(t *testing.T) {
typ := reflect.TypeOf(dataFile{})
typ := reflect.TypeFor[dataFile]()
want := make([]int, 0, typ.NumField())
for i := range typ.NumField() {
if _, ok := typ.Field(i).Tag.Lookup("avro"); ok {
Expand Down
5 changes: 1 addition & 4 deletions io/gocloud/blobfs/blob.go
Original file line number Diff line number Diff line change
Expand Up @@ -486,10 +486,7 @@ func (bfs *FileIO) DeleteFiles(ctx context.Context, paths []string) ([]string, e
}

results := make([]result, len(paths))
workers := len(paths)
if workers > deleteFilesMaxConcurrency {
workers = deleteFilesMaxConcurrency
}
workers := min(len(paths), deleteFilesMaxConcurrency)

jobs := make(chan int)
var wg sync.WaitGroup
Expand Down
2 changes: 1 addition & 1 deletion manifest_projection_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ func TestManifestEntryProjectionWhitelistCoversDataFileSchema(t *testing.T) {
schema := NewSchema(1, NestedField{
ID: 1, Name: "id", Type: PrimitiveTypes.Int64, Required: true,
})
dataFileType := reflect.TypeOf(dataFile{})
dataFileType := reflect.TypeFor[dataFile]()
dataFileFields := make(map[string]struct{}, len(dataFileAvroFieldIndexes))
for _, index := range dataFileAvroFieldIndexes {
name := dataFileType.Field(index).Tag.Get("avro")
Expand Down
12 changes: 4 additions & 8 deletions table/deferred_snapshots_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -355,13 +355,11 @@ func TestDeferredSnapshotsConcurrentMaterialization(t *testing.T) {

var wg sync.WaitGroup
for range 32 {
wg.Add(1)
go func() {
defer wg.Done()
wg.Go(func() {
require.NotNil(t, meta.CurrentSnapshot())
require.NotNil(t, meta.SnapshotByID(historicalID))
require.Len(t, meta.Snapshots(), 2)
}()
})
}
wg.Wait()
}
Expand Down Expand Up @@ -410,11 +408,9 @@ func TestDeferredSnapshotsConcurrentSingleLookupDoesNotMaterializeHistory(t *tes

var wg sync.WaitGroup
for range 32 {
wg.Add(1)
go func() {
defer wg.Done()
wg.Go(func() {
require.Equal(t, historicalID, meta.SnapshotByID(historicalID).SnapshotID)
}()
})
}
wg.Wait()

Expand Down
10 changes: 2 additions & 8 deletions table/dv/roaring_bitmap.go
Original file line number Diff line number Diff line change
Expand Up @@ -299,10 +299,7 @@ func (b *RoaringPositionBitmap) KeepMaskBytes(length int64) []byte {
if bucketBitBase >= uint64(length) {
continue
}
bucketBits := uint64(length) - bucketBitBase
if bucketBits > 1<<32 {
bucketBits = 1 << 32
}
bucketBits := min(uint64(length)-bucketBitBase, 1<<32)
if bm.CardinalityInRange(0, bucketBits) < bm.DenseSize() {
it := bm.Iterator()
for it.HasNext() {
Expand All @@ -322,10 +319,7 @@ func (b *RoaringPositionBitmap) KeepMaskBytes(length int64) []byte {
continue
}
// Cap the bucket's bit range to what fits in `length`.
bucketBits = uint64(len(dense)) * 64
if bucketBits > uint64(length)-bucketBitBase {
bucketBits = uint64(length) - bucketBitBase
}
bucketBits = min(uint64(len(dense))*64, uint64(length)-bucketBitBase)
// bucketBitBase = key << 32 is always 8-byte-aligned, so the
// BitmapWordWriter runs with offset=0 internally. The trailing-byte
// loop below relies on that alignment — PutNextTrailingByte's
Expand Down
6 changes: 2 additions & 4 deletions table/dv_scan_planning_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -230,11 +230,9 @@ func TestManifestEntries_ConcurrentMerge(t *testing.T) {
entries := newManifestEntries()
var wg sync.WaitGroup
for range manifestCount {
wg.Add(1)
go func() {
defer wg.Done()
wg.Go(func() {
assert.NoError(t, entries.merge(batch))
}()
})
}
wg.Wait()

Expand Down
6 changes: 3 additions & 3 deletions table/equality_delete_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -319,9 +319,9 @@ func schemaForEqualityFields(current *iceberg.Schema, schemas []*iceberg.Schema,
}
// Scan tasks do not retain the equality delete's sequence number, so use
// the newest schema that can resolve the file's complete equality key.
for i := len(schemas) - 1; i >= 0; i-- {
if hasAllFields(schemas[i]) {
return schemas[i]
for _, v := range slices.Backward(schemas) {
if hasAllFields(v) {
return v
}
}

Expand Down
3 changes: 1 addition & 2 deletions table/inspect_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2224,8 +2224,7 @@ func TestInspectFilesKeepCallerContextForManifestListRead(t *testing.T) {
}, nil
}

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
ctx := t.Context()
manifests, err := tbl.Inspect().Manifests(ctx)
require.NoError(t, err)
manifests.Release()
Expand Down
2 changes: 1 addition & 1 deletion table/inspect_partitions_bench_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ func benchmarkInspectPartitionFiles(
partition := make(map[int]any, benchmark.fieldCount)
for field, partitionField := range partitionFields {
if benchmark.binary {
partition[partitionField.FieldID] = []byte(fmt.Sprintf("partition-%d-field-%d", partitionID, field))
partition[partitionField.FieldID] = fmt.Appendf(nil, "partition-%d-field-%d", partitionID, field)
} else {
partition[partitionField.FieldID] = int32(partitionID + field)
}
Expand Down
4 changes: 2 additions & 2 deletions table/internal/utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -636,8 +636,8 @@ func TruncateUpperBoundBinary(val []byte, trunc int) []byte {

result := slices.Clone(val[:trunc])

for i := len(result) - 1; i >= 0; i-- {
if result[i] < 255 {
for i, v := range slices.Backward(result) {
if v < 255 {
result[i]++

return result[:i+1]
Expand Down
5 changes: 1 addition & 4 deletions table/internal/variant_shredding.go
Original file line number Diff line number Diff line change
Expand Up @@ -388,10 +388,7 @@ func decimalArrowType(info *fieldInfo) arrow.DataType {
// Always Decimal128: arrow-go's pqarrow maps it to INT32/INT64/FLBA by
// precision and cannot serialize Decimal32/Decimal64.
intDigits := max(info.maxDecimalIntDigits, 0)
prec := min(intDigits+info.maxDecimalScale, 38)
if prec < 1 {
prec = 1
}
prec := max(min(intDigits+info.maxDecimalScale, 38), 1)
scale := info.maxDecimalScale
if maxScale := 38 - intDigits; scale > maxScale {
if maxScale < 0 {
Expand Down
2 changes: 1 addition & 1 deletion table/orphan_cleanup.go
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ func flattenURIEquivalences(equivalences map[string]string) map[string]string {
continue
}

for _, value := range strings.Split(group, ",") {
for value := range strings.SplitSeq(group, ",") {
flattened[strings.TrimSpace(value)] = equivalences[group]
}
}
Expand Down
4 changes: 2 additions & 2 deletions table/positional_delete_index_bench_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,8 +97,8 @@ func positionalDeletePartitionKeyBenchmarkFiles(
partition := make(map[int]any, fieldCount)
for field, partitionField := range partitionFields {
if binaryValue {
partition[partitionField.FieldID] = []byte(fmt.Sprintf(
"partition-%02d-%02d", i%100, field))
partition[partitionField.FieldID] = fmt.Appendf(nil,
"partition-%02d-%02d", i%100, field)
} else {
partition[partitionField.FieldID] = int32(i % 100)
}
Expand Down
12 changes: 4 additions & 8 deletions table/rolling_data_writer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -315,16 +315,14 @@ func (s *RollingDataWriterTestSuite) TestConcurrentGetOrCreateCreatesOneWriter()

var wg sync.WaitGroup
for range goroutineCount {
wg.Add(1)
go func() {
defer wg.Done()
wg.Go(func() {
<-start
writer, err := factory.getOrCreateRollingDataWriter(s.ctx, "partition", nil, outputCh)
results <- struct {
writer *RollingDataWriter
err error
}{writer, err}
}()
})
}

close(start)
Expand Down Expand Up @@ -1129,16 +1127,14 @@ func (s *RollingDataWriterTestSuite) TestConcurrentAddDuringStreamErrorLeaksNoth
var atRetain, done sync.WaitGroup
atRetain.Add(senders)
for range senders {
done.Add(1)
go func() {
defer done.Done()
done.Go(func() {
record := s.buildRecord(arrSchema, 3)
// Add retains (via gatedRetainRecord.Retain, which parks the sender
// between the closed check and the enqueue) then enqueues or aborts;
// the caller always drops its own reference afterward.
_ = writer.Add(gatedRetainRecord{RecordBatch: record, retainGate: retainGate, atRetain: &atRetain})
record.Release()
}()
})
}

atRetain.Wait() // every sender is past its closed check, parked in Retain
Expand Down
4 changes: 2 additions & 2 deletions table/updates_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -730,7 +730,7 @@ func TestUnmarshalUpdatesRejectsMissingRequiredFields(t *testing.T) {
for _, tt := range tests {
t.Run(tt.action, func(t *testing.T) {
var updates Updates
data := []byte(fmt.Sprintf(`[{"action":%q}]`, tt.action))
data := fmt.Appendf(nil, `[{"action":%q}]`, tt.action)
if tt.action == UpdateSetSnapshotRef {
data = snapshotRefPayload(tt.field, false)
}
Expand Down Expand Up @@ -797,7 +797,7 @@ func TestUnmarshalUpdatesRejectsNullRequiredPayload(t *testing.T) {
for _, tt := range tests {
t.Run(tt.action+"/"+tt.field, func(t *testing.T) {
var updates Updates
data := []byte(fmt.Sprintf(`[{"action":%q,%q:null}]`, tt.action, tt.field))
data := fmt.Appendf(nil, `[{"action":%q,%q:null}]`, tt.action, tt.field)
if tt.action == UpdateSetSnapshotRef {
data = snapshotRefPayload(tt.field, true)
}
Expand Down
Loading