From 1b0d92278bf05a6625936b7995132550464df083 Mon Sep 17 00:00:00 2001 From: Neelesh Salian Date: Sun, 27 Sep 2026 13:31:39 -0700 Subject: [PATCH 1/2] refactor(catalog): move namespace-property helper into catalog/internal, drop go:linkname --- catalog/catalog.go | 51 --------------------------------------- catalog/glue/glue.go | 9 +------ catalog/hive/hive.go | 6 +---- catalog/internal/utils.go | 50 ++++++++++++++++++++++++++++++++++++++ catalog/sql/sql.go | 9 +------ 5 files changed, 53 insertions(+), 72 deletions(-) diff --git a/catalog/catalog.go b/catalog/catalog.go index b7348c304..a075b6cd3 100644 --- a/catalog/catalog.go +++ b/catalog/catalog.go @@ -32,12 +32,10 @@ import ( "errors" "fmt" "iter" - "maps" "strings" "unicode" "github.com/apache/iceberg-go" - iceinternal "github.com/apache/iceberg-go/internal" "github.com/apache/iceberg-go/table" ) @@ -357,52 +355,3 @@ func WithViewProperties(config iceberg.Properties) CreateViewOpt { cfg.Properties = config } } - -//lint:ignore U1000 this is linked to by catalogs via go:linkname but we don't want to export it -func checkForOverlap(removals []string, updates iceberg.Properties) error { - overlap := []string{} - for _, key := range removals { - if _, ok := updates[key]; ok { - overlap = append(overlap, key) - } - } - if len(overlap) > 0 { - return fmt.Errorf("conflict between removals and updates for keys: %v", overlap) - } - - return nil -} - -//lint:ignore U1000 this is linked to by catalogs via go:linkname but we don't want to export it -func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, PropertiesUpdateSummary, error) { - if err := checkForOverlap(removals, updates); err != nil { - return nil, PropertiesUpdateSummary{}, err - } - var ( - updatedProps = maps.Clone(currentProps) - removed = make([]string, 0, len(removals)) - updated = make([]string, 0, len(updates)) - ) - - for _, key := range removals { - if _, exists := updatedProps[key]; exists { - delete(updatedProps, key) - removed = append(removed, key) - } - } - - for key, value := range updates { - if updatedProps[key] != value { - updated = append(updated, key) - updatedProps[key] = value - } - } - - summary := PropertiesUpdateSummary{ - Removed: removed, - Updated: updated, - Missing: iceinternal.Difference(removals, removed), - } - - return updatedProps, summary, nil -} diff --git a/catalog/glue/glue.go b/catalog/glue/glue.go index c44aa1a7f..d33abd274 100644 --- a/catalog/glue/glue.go +++ b/catalog/glue/glue.go @@ -27,7 +27,6 @@ import ( "strconv" "strings" "time" - _ "unsafe" "github.com/apache/iceberg-go" "github.com/apache/iceberg-go/catalog" @@ -766,12 +765,6 @@ func (c *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.I return props, nil } -// avoid circular dependency while still avoiding having to export the getUpdatedPropsAndUpdateSummary function -// so that we can re-use it in the catalog implementations without duplicating the code. - -//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary -func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error) - // UpdateNamespaceProperties updates the properties of an Iceberg namespace in the Glue catalog. // The removals list contains the keys to remove, and the updates map contains the keys and values to update. func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, @@ -782,7 +775,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table return catalog.PropertiesUpdateSummary{}, err } - updatedProperties, propertiesUpdateSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates) + updatedProperties, propertiesUpdateSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates) if err != nil { return catalog.PropertiesUpdateSummary{}, err } diff --git a/catalog/hive/hive.go b/catalog/hive/hive.go index 604d5b1d5..af4e206fb 100644 --- a/catalog/hive/hive.go +++ b/catalog/hive/hive.go @@ -25,7 +25,6 @@ import ( "log" "maps" "strings" - _ "unsafe" "github.com/apache/iceberg-go" "github.com/apache/iceberg-go/catalog" @@ -934,9 +933,6 @@ func (c *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.I return props, nil } -//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary -func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error) - // UpdateNamespaceProperties updates the properties for a namespace. func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, removals []string, updates iceberg.Properties, @@ -946,7 +942,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table return catalog.PropertiesUpdateSummary{}, err } - updatedProperties, propertiesUpdateSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates) + updatedProperties, propertiesUpdateSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates) if err != nil { return catalog.PropertiesUpdateSummary{}, err } diff --git a/catalog/internal/utils.go b/catalog/internal/utils.go index fbb088e83..41e9a03f6 100644 --- a/catalog/internal/utils.go +++ b/catalog/internal/utils.go @@ -324,3 +324,53 @@ func UpdateAndStageTable(ctx context.Context, catprops iceberg.Properties, curre ), }, nil } + +func checkForOverlap(removals []string, updates iceberg.Properties) error { + overlap := []string{} + for _, key := range removals { + if _, ok := updates[key]; ok { + overlap = append(overlap, key) + } + } + if len(overlap) > 0 { + return fmt.Errorf("conflict between removals and updates for keys: %v", overlap) + } + + return nil +} + +// GetUpdatedPropsAndUpdateSummary applies removals and updates to currentProps +// and returns the updated properties alongside a summary of the changes. It is +// shared by the catalog backend implementations. +func GetUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error) { + if err := checkForOverlap(removals, updates); err != nil { + return nil, catalog.PropertiesUpdateSummary{}, err + } + var ( + updatedProps = maps.Clone(currentProps) + removed = make([]string, 0, len(removals)) + updated = make([]string, 0, len(updates)) + ) + + for _, key := range removals { + if _, exists := updatedProps[key]; exists { + delete(updatedProps, key) + removed = append(removed, key) + } + } + + for key, value := range updates { + if updatedProps[key] != value { + updated = append(updated, key) + updatedProps[key] = value + } + } + + summary := catalog.PropertiesUpdateSummary{ + Removed: removed, + Updated: updated, + Missing: internal.Difference(removals, removed), + } + + return updatedProps, summary, nil +} diff --git a/catalog/sql/sql.go b/catalog/sql/sql.go index 39c2d20c4..61370d1c6 100644 --- a/catalog/sql/sql.go +++ b/catalog/sql/sql.go @@ -30,7 +30,6 @@ import ( "strings" "sync" "time" - _ "unsafe" "github.com/apache/iceberg-go" "github.com/apache/iceberg-go/catalog" @@ -1518,12 +1517,6 @@ func (c *Catalog) ListNamespaces(ctx context.Context, parent table.Identifier) ( return ret, nil } -// avoid circular dependency while still avoiding having to export the getUpdatedPropsAndUpdateSummary function -// so that we can re-use it in the catalog implementations without duplicating the code. - -//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary -func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error) - func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, removals []string, updates iceberg.Properties) (catalog.PropertiesUpdateSummary, error) { var summary catalog.PropertiesUpdateSummary if err := checkValidNamespace(namespace); err != nil { @@ -1555,7 +1548,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table currentProps[prop.PropertyKey] = prop.PropertyValue.String } - _, nextSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates) + _, nextSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates) if err != nil { return err } From 91abe5ee1ddde4f141b12a59fb4459d1b238f9e0 Mon Sep 17 00:00:00 2001 From: Neelesh Salian Date: Mon, 28 Sep 2026 13:09:35 -0700 Subject: [PATCH 2/2] PR fixes --- catalog/internal/utils.go | 10 +-- catalog/internal/utils_test.go | 137 +++++++++++++++++++++++++++++++++ 2 files changed, 142 insertions(+), 5 deletions(-) create mode 100644 catalog/internal/utils_test.go diff --git a/catalog/internal/utils.go b/catalog/internal/utils.go index 41e9a03f6..203ac680d 100644 --- a/catalog/internal/utils.go +++ b/catalog/internal/utils.go @@ -33,7 +33,7 @@ import ( "github.com/apache/iceberg-go" "github.com/apache/iceberg-go/catalog" - "github.com/apache/iceberg-go/internal" + iceinternal "github.com/apache/iceberg-go/internal" icebergio "github.com/apache/iceberg-go/io" "github.com/apache/iceberg-go/table" "github.com/google/uuid" @@ -52,21 +52,21 @@ func WriteTableMetadata(metadata table.Metadata, fs icebergio.WriteFileIO, loc s if err != nil { return err } - defer internal.CheckedClose(out, &err) + defer iceinternal.CheckedClose(out, &err) var writer io.Writer = out switch compression { case table.MetadataCompressionCodecGzip: gzw := gzip.NewWriter(out) writer = gzw - defer internal.CheckedClose(gzw, &err) + defer iceinternal.CheckedClose(gzw, &err) case table.MetadataCompressionCodecZstd: enc, zErr := zstd.NewWriter(out) if zErr != nil { return zErr } writer = enc - defer internal.CheckedClose(enc, &err) + defer iceinternal.CheckedClose(enc, &err) } err = json.NewEncoder(writer).Encode(metadata) @@ -369,7 +369,7 @@ func GetUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals [ summary := catalog.PropertiesUpdateSummary{ Removed: removed, Updated: updated, - Missing: internal.Difference(removals, removed), + Missing: iceinternal.Difference(removals, removed), } return updatedProps, summary, nil diff --git a/catalog/internal/utils_test.go b/catalog/internal/utils_test.go new file mode 100644 index 000000000..319e50236 --- /dev/null +++ b/catalog/internal/utils_test.go @@ -0,0 +1,137 @@ +// 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 internal + +import ( + "maps" + "testing" + + "github.com/apache/iceberg-go" + "github.com/apache/iceberg-go/catalog" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestGetUpdatedPropsAndUpdateSummary(t *testing.T) { + tests := []struct { + name string + current iceberg.Properties + removals []string + updates iceberg.Properties + wantErrContains string + wantProps iceberg.Properties + wantRemoved []string + wantUpdated []string + wantMissing []string + }{ + { + name: "removal and update of the same key conflicts", + current: iceberg.Properties{"a": "1"}, + removals: []string{"a"}, + updates: iceberg.Properties{"a": "2"}, + wantErrContains: "conflict between removals and updates for keys: [a]", + }, + { + name: "existing key is removed", + current: iceberg.Properties{"a": "1", "b": "2"}, + removals: []string{"a"}, + wantProps: iceberg.Properties{"b": "2"}, + wantRemoved: []string{"a"}, + wantUpdated: []string{}, + wantMissing: []string{}, + }, + { + name: "removal of absent key is reported as missing", + current: iceberg.Properties{"a": "1"}, + removals: []string{"ghost"}, + wantProps: iceberg.Properties{"a": "1"}, + wantRemoved: []string{}, + wantUpdated: []string{}, + wantMissing: []string{"ghost"}, + }, + { + name: "existing key gets a new value", + current: iceberg.Properties{"a": "old"}, + updates: iceberg.Properties{"a": "new"}, + wantProps: iceberg.Properties{"a": "new"}, + wantRemoved: []string{}, + wantUpdated: []string{"a"}, + wantMissing: []string{}, + }, + { + name: "new key is added", + current: iceberg.Properties{"a": "1"}, + updates: iceberg.Properties{"b": "2"}, + wantProps: iceberg.Properties{"a": "1", "b": "2"}, + wantRemoved: []string{}, + wantUpdated: []string{"b"}, + wantMissing: []string{}, + }, + { + name: "update with unchanged value is not reported", + current: iceberg.Properties{"a": "1"}, + updates: iceberg.Properties{"a": "1"}, + wantProps: iceberg.Properties{"a": "1"}, + wantRemoved: []string{}, + wantUpdated: []string{}, + wantMissing: []string{}, + }, + { + name: "no removals or updates is a no-op", + current: iceberg.Properties{"a": "1"}, + wantProps: iceberg.Properties{"a": "1"}, + wantRemoved: []string{}, + wantUpdated: []string{}, + wantMissing: []string{}, + }, + { + name: "add, update, remove and missing in one call", + current: iceberg.Properties{"keep": "v", "drop": "v", "upd": "old"}, + removals: []string{"drop", "ghost"}, + updates: iceberg.Properties{"upd": "new", "add": "v", "keep": "v"}, + wantProps: iceberg.Properties{"keep": "v", "upd": "new", "add": "v"}, + wantRemoved: []string{"drop"}, + wantUpdated: []string{"upd", "add"}, + wantMissing: []string{"ghost"}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + before := maps.Clone(tt.current) + + props, summary, err := GetUpdatedPropsAndUpdateSummary(tt.current, tt.removals, tt.updates) + + assert.Equal(t, before, tt.current, "currentProps must not be mutated") + + if tt.wantErrContains != "" { + require.ErrorContains(t, err, tt.wantErrContains) + assert.Nil(t, props) + assert.Equal(t, catalog.PropertiesUpdateSummary{}, summary) + + return + } + + require.NoError(t, err) + assert.Equal(t, tt.wantProps, props) + assert.ElementsMatch(t, tt.wantRemoved, summary.Removed) + assert.ElementsMatch(t, tt.wantUpdated, summary.Updated) + assert.ElementsMatch(t, tt.wantMissing, summary.Missing) + }) + } +}