diff --git a/catalog/rest/rest.go b/catalog/rest/rest.go index 2f5af16df..1e78bb5d1 100644 --- a/catalog/rest/rest.go +++ b/catalog/rest/rest.go @@ -1244,6 +1244,7 @@ func (r *Catalog) tableFromResponse( config iceberg.Properties, scanPlanningConfig iceberg.Properties, credsVended bool, + labels *iceberg.Labels, ) (*table.Table, error) { var fsF func(context.Context) (iceio.IO, error) if credsVended { @@ -1302,6 +1303,7 @@ func (r *Catalog) tableFromResponse( r, table.WithMetricsReporter(reporter), table.WithScanPlanningIOProperties(scanPlanningConfig), + table.WithLabels(labels), ), nil } @@ -1541,7 +1543,7 @@ func (r *Catalog) CreateTable(ctx context.Context, identifier table.Identifier, credsVended := len(ret.StorageCredentials) > 0 maps.Copy(config, resolveStorageCredentials(ret.StorageCredentials, ret.MetadataLoc)) - return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended) + return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels) } // commitStagedCreate performs the second phase of a staged table @@ -1769,7 +1771,7 @@ func (r *Catalog) RegisterTable(ctx context.Context, identifier table.Identifier credsVended := len(ret.StorageCredentials) > 0 maps.Copy(config, resolveStorageCredentials(ret.StorageCredentials, ret.MetadataLoc)) - return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended) + return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels) } // LoadTable loads a table from the catalog. It implements [catalog.Catalog]. @@ -1814,7 +1816,7 @@ func (r *Catalog) loadTableWithMode(ctx context.Context, identifier table.Identi credsVended := len(ret.StorageCredentials) > 0 maps.Copy(config, resolveStorageCredentials(ret.StorageCredentials, ret.MetadataLoc)) - return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended) + return r.tableFromResponse(ctx, identifier, ret.Metadata, ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels) } func (r *Catalog) UpdateTable(ctx context.Context, ident table.Identifier, requirements []table.Requirement, updates []table.Update) (*table.Table, error) { @@ -1862,7 +1864,8 @@ func (r *Catalog) UpdateTable(ctx context.Context, ident table.Identifier, requi config := maps.Clone(r.props) maps.Copy(config, metadata.Properties()) - return r.tableFromResponse(ctx, ident, metadata, ret.MetadataLoc, config, config, false) + // A commit response carries no labels (they are load-time enrichment). + return r.tableFromResponse(ctx, ident, metadata, ret.MetadataLoc, config, config, false, nil) } func (r *Catalog) DropTable(ctx context.Context, identifier table.Identifier) error { @@ -2339,6 +2342,7 @@ type viewResponse struct { MetadataLoc string `json:"metadata-location"` RawMetadata json.RawMessage `json:"metadata"` Config iceberg.Properties `json:"config"` + Labels *iceberg.Labels `json:"labels,omitempty"` Metadata view.Metadata `json:"-"` } @@ -2414,7 +2418,7 @@ func (r *Catalog) CreateView(ctx context.Context, identifier table.Identifier, v return nil, err } - return view.New(identifier, ret.Metadata, ret.MetadataLoc), nil + return view.New(identifier, ret.Metadata, ret.MetadataLoc, view.WithLabels(ret.Labels)), nil } // UpdateView updates a view in the catalog. @@ -2449,7 +2453,7 @@ func (r *Catalog) UpdateView(ctx context.Context, ident table.Identifier, requir return nil, err } - return view.New(ident, ret.Metadata, ret.MetadataLoc), nil + return view.New(ident, ret.Metadata, ret.MetadataLoc, view.WithLabels(ret.Labels)), nil } // loadViewResponse contains the response from loading a view @@ -2497,7 +2501,7 @@ func (r *Catalog) RegisterView(ctx context.Context, identifier table.Identifier, return nil, fmt.Errorf("failed to parse view metadata: %w", err) } - return view.New(identifier, metadata, rsp.MetadataLoc), nil + return view.New(identifier, metadata, rsp.MetadataLoc, view.WithLabels(rsp.Labels)), nil } // LoadView loads a view from the catalog. @@ -2529,7 +2533,7 @@ func (r *Catalog) LoadView(ctx context.Context, identifier table.Identifier) (*v return nil, fmt.Errorf("failed to parse view metadata: %w", err) } - return view.New(identifier, metadata, rsp.MetadataLoc), nil + return view.New(identifier, metadata, rsp.MetadataLoc, view.WithLabels(rsp.Labels)), nil } func (r *Catalog) RenameView(ctx context.Context, from, to table.Identifier) (*view.View, error) { diff --git a/catalog/rest/rest_test.go b/catalog/rest/rest_test.go index 9b2749960..dc4f6653e 100644 --- a/catalog/rest/rest_test.go +++ b/catalog/rest/rest_test.go @@ -1633,6 +1633,98 @@ func (r *RestCatalogSuite) TestLoadTable200() { }, }, })) + + r.Nil(tbl.Labels(), "no labels in response should surface as nil") +} + +func (r *RestCatalogSuite) TestLoadTableLabels() { + r.mux.HandleFunc("/v1/namespaces/fokko/tables/table", func(w http.ResponseWriter, req *http.Request) { + r.Require().Equal(http.MethodGet, req.Method) + w.Write([]byte(`{ + "metadata-location": "s3://warehouse/database/table/metadata/00001.metadata.json", + "metadata": { + "format-version": 1, + "table-uuid": "b55d9dda-6561-423a-8bfc-787980ce421f", + "location": "s3://warehouse/database/table", + "last-updated-ms": 1646787054459, + "last-column-id": 2, + "schema": {"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}, + "current-schema-id": 0, + "schemas": [{"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}], + "partition-spec": [], + "default-spec-id": 0, + "partition-specs": [{"spec-id":0,"fields":[]}], + "last-partition-id": 999, + "default-sort-order-id": 0, + "sort-orders": [{"order-id":0,"fields":[]}], + "properties": {} + }, + "config": {}, + "labels": { + "object-labels": {"owner": "data-eng", "cost-center": "42"}, + "fields": [ + {"field-id": 2, "labels": {"classification": "pii"}}, + {"field-id": 99, "labels": {"classification": "dropped"}} + ] + } + }`)) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL, rest.WithOAuthToken(TestToken)) + r.Require().NoError(err) + + tbl, err := cat.LoadTable(context.Background(), catalog.ToIdentifier("fokko", "table")) + r.Require().NoError(err) + + labels := tbl.Labels() + r.Require().NotNil(labels) + r.False(labels.IsEmpty()) + r.Equal(iceberg.Properties{"owner": "data-eng", "cost-center": "42"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "pii"}, labels.Field(2)) + r.Equal(iceberg.Properties{"classification": "dropped"}, labels.Field(99)) // dropped column still resolvable + r.Nil(labels.Field(1)) // field present in schema but unlabeled +} + +func (r *RestCatalogSuite) TestCreateTableLabels() { + r.mux.HandleFunc("/v1/namespaces/fokko/tables", func(w http.ResponseWriter, req *http.Request) { + r.Require().Equal(http.MethodPost, req.Method) + w.Write([]byte(`{ + "metadata-location": "s3://warehouse/database/table/metadata/00001.metadata.json", + "metadata": { + "format-version": 1, + "table-uuid": "b55d9dda-6561-423a-8bfc-787980ce421f", + "location": "s3://warehouse/database/table", + "last-updated-ms": 1646787054459, + "last-column-id": 2, + "schema": {"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}, + "current-schema-id": 0, + "schemas": [{"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}], + "partition-spec": [], + "default-spec-id": 0, + "partition-specs": [{"spec-id":0,"fields":[]}], + "last-partition-id": 999, + "default-sort-order-id": 0, + "sort-orders": [{"order-id":0,"fields":[]}], + "properties": {} + }, + "config": {}, + "labels": { + "object-labels": {"owner": "analytics"}, + "fields": [{"field-id": 2, "labels": {"classification": "pii"}}] + } + }`)) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL, rest.WithOAuthToken(TestToken)) + r.Require().NoError(err) + + tbl, err := cat.CreateTable(context.Background(), catalog.ToIdentifier("fokko", "fokko2"), tableSchemaSimple) + r.Require().NoError(err) + + labels := tbl.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "analytics"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "pii"}, labels.Field(2)) } func (r *RestCatalogSuite) TestLoadTableWithSnapshotModeRefs() { @@ -2041,6 +2133,46 @@ func (r *RestCatalogSuite) TestRegisterTable200() { r.Equal("bryan", tbl.Metadata().Properties()["owner"]) } +func (r *RestCatalogSuite) TestRegisterTableLabels() { + r.mux.HandleFunc("/v1/namespaces/fokko/register", func(w http.ResponseWriter, req *http.Request) { + r.Require().Equal(http.MethodPost, req.Method) + w.Write([]byte(`{ + "metadata-location": "s3://warehouse/database/table/metadata/00001.metadata.json", + "metadata": { + "format-version": 1, + "table-uuid": "b55d9dda-6561-423a-8bfc-787980ce421f", + "location": "s3://warehouse/database/table", + "last-updated-ms": 1646787054459, + "last-column-id": 2, + "schema": {"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}, + "current-schema-id": 0, + "schemas": [{"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"int"},{"id":2,"name":"data","required":false,"type":"string"}]}], + "partition-spec": [], + "default-spec-id": 0, + "partition-specs": [{"spec-id":0,"fields":[]}], + "last-partition-id": 999, + "default-sort-order-id": 0, + "sort-orders": [{"order-id":0,"fields":[]}], + "properties": {} + }, + "config": {}, + "labels": {"object-labels": {"owner": "data-eng"}, "fields": [{"field-id": 2, "labels": {"classification": "pii"}}]} + }`)) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL, rest.WithOAuthToken(TestToken)) + r.Require().NoError(err) + + tbl, err := cat.RegisterTable(context.Background(), catalog.ToIdentifier("fokko", "fokko2"), + "s3://warehouse/database/table/metadata/00001.metadata.json") + r.Require().NoError(err) + + labels := tbl.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "data-eng"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "pii"}, labels.Field(2)) +} + func (r *RestCatalogSuite) TestRegisterTable404() { r.mux.HandleFunc("/v1/namespaces/nonexistent/register", func(w http.ResponseWriter, req *http.Request) { r.Require().Equal(http.MethodPost, req.Method) @@ -2568,6 +2700,55 @@ func (r *RestCatalogSuite) TestLoadView200() { r.Len(currentVersion.Representations, 1) r.Equal("sql", currentVersion.Representations[0].Type) r.Equal("spark", currentVersion.Representations[0].Dialect) + + r.Nil(v.Labels(), "no labels in response should surface as nil") +} + +func (r *RestCatalogSuite) TestLoadViewLabels() { + r.mux.HandleFunc("/v1/namespaces/fokko/views/myview", func(w http.ResponseWriter, req *http.Request) { + r.Require().Equal(http.MethodGet, req.Method) + w.Write([]byte(`{ + "metadata-location": "s3://bucket/warehouse/default.db/event_agg/metadata/00001.metadata.json", + "metadata": { + "view-uuid": "fa6506c3-7681-40c8-86dc-e36561f83385", + "format-version": 1, + "location": "s3://bucket/warehouse/default.db/event_agg", + "current-version-id": 1, + "properties": {}, + "versions": [{ + "version-id": 1, + "timestamp-ms": 1573518431292, + "schema-id": 1, + "default-catalog": "prod", + "default-namespace": ["default"], + "summary": {"engine-name": "Spark"}, + "representations": [{"type": "sql", "sql": "SELECT 1", "dialect": "spark"}] + }], + "schemas": [{"schema-id": 1, "type": "struct", "fields": [ + {"id": 1, "name": "event_count", "required": false, "type": "int"}, + {"id": 2, "name": "event_date", "required": false, "type": "date"} + ]}], + "version-log": [{"timestamp-ms": 1573518431292, "version-id": 1}] + }, + "config": {}, + "labels": { + "object-labels": {"owner": "analytics"}, + "fields": [{"field-id": 2, "labels": {"classification": "internal"}}] + } + }`)) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL, rest.WithOAuthToken(TestToken)) + r.Require().NoError(err) + + v, err := cat.LoadView(context.Background(), catalog.ToIdentifier("fokko", "myview")) + r.Require().NoError(err) + + labels := v.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "analytics"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "internal"}, labels.Field(2)) + r.Nil(labels.Field(1)) // field present in schema but unlabeled } func (r *RestCatalogSuite) TestLoadView404() { @@ -2953,6 +3134,87 @@ func (r *RestCatalogSuite) TestRegisterView200() { r.Equal(exampleViewSQL, v.Metadata().CurrentVersion().Representations[0].Sql) } +func (r *RestCatalogSuite) TestRegisterViewLabels() { + const ( + ns = "fokko" + viewName = "myview" + metadataLoc = "s3://bucket/warehouse/fokko.db/myview/metadata/00001.metadata.json" + ) + + r.mux.HandleFunc("/v1/namespaces/"+ns+"/register-view", func(w http.ResponseWriter, req *http.Request) { + r.Require().Equal(http.MethodPost, req.Method) + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"metadata-location": %q, "metadata": %s, "config": {}, "labels": {"object-labels": {"owner": "analytics"}}}`, + metadataLoc, exampleViewMetadataJSON) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL, rest.WithOAuthToken(TestToken)) + r.Require().NoError(err) + + v, err := cat.RegisterView(context.Background(), table.Identifier{ns, viewName}, metadataLoc) + r.Require().NoError(err) + + labels := v.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "analytics"}, labels.Object()) +} + +func (r *RestCatalogSuite) TestCreateViewLabels() { + ns := "ns" + viewName := "view" + metadataLoc := "s3://bucket/warehouse/ns.db/view/metadata/00001.metadata.json" + identifier := table.Identifier{ns, viewName} + schema := iceberg.NewSchemaWithIdentifiers(0, []int{1}, iceberg.NestedField{ + ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int32, Required: true, + }) + reprs := []view.Representation{view.NewRepresentation(exampleViewSQL, "default")} + version, err := view.NewVersion(1, 0, reprs, table.Identifier{ns}, + view.WithDefaultViewCatalog("default-catalog"), view.WithTimestampMS(0)) + r.Require().NoError(err) + + r.mux.HandleFunc("/v1/namespaces/"+ns+"/views", func(w http.ResponseWriter, req *http.Request) { + r.Equal(http.MethodPost, req.Method) + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"metadata-location": %q, "metadata": %s, "config": {}, "labels": {"object-labels": {"owner": "analytics"}, "fields": [{"field-id": 1, "labels": {"classification": "internal"}}]}}`, + metadataLoc, exampleViewMetadataJSON) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL) + r.Require().NoError(err) + + v, err := cat.CreateView(context.Background(), identifier, version, schema) + r.Require().NoError(err) + + labels := v.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "analytics"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "internal"}, labels.Field(1)) +} + +func (r *RestCatalogSuite) TestUpdateViewLabels() { + ns := "ns" + viewName := "view" + metadataLoc := "s3://bucket/warehouse/ns.db/view/metadata/00002.metadata.json" + + r.mux.HandleFunc("/v1/namespaces/"+ns+"/views/"+viewName, func(w http.ResponseWriter, req *http.Request) { + r.Equal(http.MethodPost, req.Method) + w.Header().Set("Content-Type", "application/json") + fmt.Fprintf(w, `{"metadata-location": %q, "metadata": %s, "config": {}, "labels": {"object-labels": {"owner": "analytics"}, "fields": [{"field-id": 1, "labels": {"classification": "internal"}}]}}`, + metadataLoc, exampleViewMetadataJSON) + }) + + cat, err := rest.NewCatalog(context.Background(), "rest", r.srv.URL) + r.Require().NoError(err) + + v, err := cat.UpdateView(context.Background(), table.Identifier{ns, viewName}, nil, nil) + r.Require().NoError(err) + + labels := v.Labels() + r.Require().NotNil(labels) + r.Equal(iceberg.Properties{"owner": "analytics"}, labels.Object()) + r.Equal(iceberg.Properties{"classification": "internal"}, labels.Field(1)) +} + func (r *RestCatalogSuite) TestRegisterView404() { const ( ns = "nonexistent" diff --git a/table/labels_wiring_test.go b/table/labels_wiring_test.go new file mode 100644 index 000000000..535867e08 --- /dev/null +++ b/table/labels_wiring_test.go @@ -0,0 +1,151 @@ +// 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" + "path/filepath" + "testing" + + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// labelStubCatalog hands every LoadTable a table carrying the configured labels, +// so Refresh can be checked to re-hydrate them. +type labelStubCatalog struct { + labels *iceberg.Labels +} + +func (c *labelStubCatalog) LoadTable(_ context.Context, ident Identifier) (*Table, error) { + return New(ident, nil, "", + func(context.Context) (iceio.IO, error) { return iceio.LocalFS{}, nil }, c, + WithLabels(c.labels)), nil +} + +func (c *labelStubCatalog) CommitTable(context.Context, Identifier, []Requirement, []Update) (Metadata, string, error) { + return nil, "", nil +} + +// committingLabelCatalog applies updates on CommitTable so a transaction Commit +// can be exercised end to end against stubbed metadata. +type committingLabelCatalog struct { + metadata Metadata +} + +func (c *committingLabelCatalog) LoadTable(_ context.Context, ident Identifier) (*Table, error) { + return New(ident, c.metadata, "", + func(context.Context) (iceio.IO, error) { return iceio.LocalFS{}, nil }, c), nil +} + +func (c *committingLabelCatalog) CommitTable(_ context.Context, _ Identifier, _ []Requirement, updates []Update) (Metadata, string, error) { + meta, err := UpdateTableMetadata(c.metadata, updates, "") + if err != nil { + return nil, "", err + } + c.metadata = meta + + return meta, "", nil +} + +// TestRefreshRehydratesLabels pins that Refresh adopts the labels from the +// reloaded table. +func TestRefreshRehydratesLabels(t *testing.T) { + labels := &iceberg.Labels{ObjectLabels: iceberg.Properties{"owner": "analytics"}} + cat := &labelStubCatalog{labels: labels} + tbl := New(Identifier{"db", "labels_refresh"}, nil, "", + func(context.Context) (iceio.IO, error) { return iceio.LocalFS{}, nil }, cat) + + require.Nil(t, tbl.Labels()) + require.NoError(t, tbl.Refresh(context.Background())) + require.NotNil(t, tbl.Labels()) + assert.Equal(t, iceberg.Properties{"owner": "analytics"}, tbl.Labels().Object()) +} + +// TestCommitPreservesLabels pins the commit path: doCommit rebuilds the table via +// New(...), so without WithLabels the returned table would carry Labels() == nil +// until a Refresh. The *Table returned by Commit must keep them. +func TestCommitPreservesLabels(t *testing.T) { + schema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + ) + loc := filepath.ToSlash(t.TempDir()) + meta, err := NewMetadata(schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, loc, + iceberg.Properties{PropertyFormatVersion: "2"}) + require.NoError(t, err) + + labels := &iceberg.Labels{ + ObjectLabels: iceberg.Properties{"owner": "analytics"}, + Fields: []iceberg.FieldLabel{{FieldID: 1, Labels: iceberg.Properties{"classification": "internal"}}}, + } + cat := &committingLabelCatalog{metadata: meta} + tbl := New(Identifier{"db", "commit_labels"}, meta, loc+"/metadata/v1.metadata.json", + func(context.Context) (iceio.IO, error) { return iceio.LocalFS{}, nil }, cat, + WithLabels(labels)) + + txn := tbl.NewTransaction() + require.NoError(t, txn.SetProperties(iceberg.Properties{"k": "v"})) + committed, err := txn.Commit(context.Background()) + require.NoError(t, err) + + require.NotNil(t, committed.Labels(), "commit must carry labels onto the returned table") + assert.Equal(t, iceberg.Properties{"owner": "analytics"}, committed.Labels().Object()) + assert.Equal(t, iceberg.Properties{"classification": "internal"}, committed.Labels().Field(1)) +} + +// TestStagedTablePreservesLabels pins that Transaction.StagedTable() forwards the +// transaction table's labels; it rebuilds the table via New(...) too. +func TestStagedTablePreservesLabels(t *testing.T) { + schema := iceberg.NewSchema(1, iceberg.NestedField{ + ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true, + }) + meta, err := NewMetadata(schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, + "mem://default/staged", iceberg.Properties{PropertyFormatVersion: "2"}) + require.NoError(t, err) + + labels := &iceberg.Labels{ObjectLabels: iceberg.Properties{"owner": "analytics"}} + tbl := New(Identifier{"default", "staged"}, meta, "", + func(context.Context) (iceio.IO, error) { return iceio.NewMemFS(), nil }, nil, + WithLabels(labels)) + + staged, err := tbl.NewTransaction().StagedTable() + require.NoError(t, err) + require.NotNil(t, staged.Labels(), "StagedTable must forward the transaction table's labels") + assert.Equal(t, iceberg.Properties{"owner": "analytics"}, staged.Labels().Object()) +} + +// TestEqualsIgnoresLabels pins that labels are transient enrichment, excluded +// from Table.Equals. +func TestEqualsIgnoresLabels(t *testing.T) { + schema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + ) + meta, err := NewMetadata(schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, "mem://db/tbl", + iceberg.Properties{PropertyFormatVersion: "2"}) + require.NoError(t, err) + + fsF := func(context.Context) (iceio.IO, error) { return iceio.NewMemFS(), nil } + labeled := New(Identifier{"db", "tbl"}, meta, "loc", fsF, nil, + WithLabels(&iceberg.Labels{ObjectLabels: iceberg.Properties{"owner": "analytics"}})) + plain := New(Identifier{"db", "tbl"}, meta, "loc", fsF, nil) + + assert.True(t, labeled.Equals(*plain), "labels must not affect Equals") + assert.True(t, plain.Equals(*labeled)) +} diff --git a/table/table.go b/table/table.go index 39ffa8649..6771c6ad8 100644 --- a/table/table.go +++ b/table/table.go @@ -117,6 +117,9 @@ type Table struct { // explicit NopReporter opt-out) from the construction-time default, so // Refresh knows whether it may overwrite reporter with the catalog default. reporterSet bool + // labels is transient catalog enrichment from the load response; nil when + // the catalog returned none. It is excluded from Equals. + labels *iceberg.Labels } func (t Table) Equals(other Table) bool { @@ -134,6 +137,12 @@ func (t Table) Spec() iceberg.PartitionSpec { return t.metadata func (t Table) SortOrder() SortOrder { return t.metadata.SortOrder() } func (t Table) Properties() iceberg.Properties { return t.metadata.Properties() } +// Labels returns the catalog-provided labels from the load response, or nil if +// the catalog returned none. Labels are transient enrichment, not table state. +// The returned pointer aliases the table's labels and must be treated as +// read-only; mutating it affects the shared table value. +func (t Table) Labels() *iceberg.Labels { return t.labels } + // MetricsReporter returns the table's metrics reporter, never nil. func (t Table) MetricsReporter() metrics.Reporter { if t.reporter == nil { @@ -239,6 +248,7 @@ func (t *Table) Refresh(ctx context.Context) error { t.manifestCache = newSnapshotManifestCacheForMetadata(fresh.metadata) t.planner = fresh.planner t.scanPlanningIOProps = maps.Clone(fresh.scanPlanningIOProps) + t.labels = fresh.labels // Only inherit the catalog-derived reporter when the caller hasn't set one // of their own. Refresh runs inside commit retry loops, so unconditionally // copying fresh.reporter would silently revert a WithMetricsReporter-injected @@ -825,6 +835,7 @@ func (t Table) doCommit(ctx context.Context, updates []Update, reqs []Requiremen t.cat, withReporterState(t.reporter, t.reporterSet), WithScanPlanningIOProperties(t.scanPlanningIOProps), + WithLabels(t.labels), ), nil } @@ -1400,6 +1411,18 @@ func WithScanPlanningIOProperties(props iceberg.Properties) Option { } } +// WithLabels attaches catalog-provided labels from a load response to the +// table. A nil value is ignored, leaving the table's labels nil. +func WithLabels(l *iceberg.Labels) Option { + if l == nil { + return noopTableOption + } + + return func(t *Table) { + t.labels = l + } +} + // withReporterState copies both the reporter and the reporterSet flag verbatim. // Unlike WithMetricsReporter it does not force reporterSet true, so a table // riding the catalog default (reporterSet == false) stays defaulted across a diff --git a/table/transaction.go b/table/transaction.go index 22687295f..b83acc33d 100644 --- a/table/transaction.go +++ b/table/transaction.go @@ -3327,6 +3327,7 @@ func (t *Transaction) StagedTable() (*StagedTable, error) { t.tbl.cat, withReporterState(t.tbl.reporter, t.tbl.reporterSet), WithScanPlanningIOProperties(t.tbl.scanPlanningIOProps), + WithLabels(t.tbl.labels), ), }, nil } diff --git a/view/labels_wiring_test.go b/view/labels_wiring_test.go new file mode 100644 index 000000000..7c72a40d2 --- /dev/null +++ b/view/labels_wiring_test.go @@ -0,0 +1,43 @@ +// 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 view + +import ( + "testing" + + "github.com/apache/iceberg-go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestViewEqualsIgnoresLabels pins that labels are transient and do not +// participate in View.Equals; two otherwise-identical views compare equal +// regardless of labels. +func TestViewEqualsIgnoresLabels(t *testing.T) { + meta, err := ParseMetadataString(exampleViewJSON) + require.NoError(t, err) + + ident := []string{"foo"} + loc := "s3://bucket/test/location/uuid.metadata.json" + withLabels := New(ident, meta, loc, + WithLabels(&iceberg.Labels{ObjectLabels: iceberg.Properties{"owner": "analytics"}})) + without := New(ident, meta, loc) + + assert.True(t, withLabels.Equals(*without), "labels are transient and must not affect Equals") + assert.True(t, without.Equals(*withLabels)) +} diff --git a/view/view.go b/view/view.go index a12be3938..3a97e0ab3 100644 --- a/view/view.go +++ b/view/view.go @@ -36,6 +36,9 @@ type View struct { identifier table.Identifier metadata Metadata metadataLocation string + // labels is transient catalog enrichment from the load response; nil when + // the catalog returned none. It is excluded from Equals. + labels *iceberg.Labels } func (t View) Equals(other View) bool { @@ -54,12 +57,42 @@ func (t View) Location() string { return t.metadata.Location() } func (t View) Versions() []*Version { return t.metadata.Versions() } func (t View) Schemas() map[int]*iceberg.Schema { return t.metadata.SchemasByID() } -func New(ident table.Identifier, meta Metadata, metadataLocation string) *View { - return &View{ +// Labels returns the catalog-provided labels from the load response, or nil if +// the catalog returned none. Labels are transient enrichment, not view state. +// The returned pointer aliases the view's labels and must be treated as +// read-only; mutating it affects the shared view value. +func (t View) Labels() *iceberg.Labels { return t.labels } + +// Option configures a [View] at construction. +type Option func(*View) + +// noopViewOption is the shared no-op [Option], returned when an option has +// nothing to apply (e.g. WithLabels(nil)). It mirrors noopTableOption. +func noopViewOption(*View) {} + +// WithLabels attaches catalog-provided labels from a load response to the view. +// A nil value is ignored, leaving the view's labels nil. +func WithLabels(l *iceberg.Labels) Option { + if l == nil { + return noopViewOption + } + + return func(v *View) { + v.labels = l + } +} + +func New(ident table.Identifier, meta Metadata, metadataLocation string, opts ...Option) *View { + v := &View{ identifier: slices.Clone(ident), metadata: meta, metadataLocation: metadataLocation, } + for _, opt := range opts { + opt(v) + } + + return v } func NewFromLocation(