From c4781f57d462556146e0b94a3e4ffc5fdacb13cb Mon Sep 17 00:00:00 2001 From: laskoviymishka Date: Tue, 15 Sep 2026 15:52:34 +0200 Subject: [PATCH] feat(rest): expose catalog labels on the loaded table and view Thread the labels decoded from load responses onto the returned Table and View. Table gains a WithLabels option and a Labels() accessor; View gains a functional-options constructor with the same. LoadTable, RegisterTable, CreateTable, LoadView, RegisterView, CreateView, and UpdateView forward the response labels (viewResponse now decodes them too). The commit path keeps labels on the returned table too: doCommit and Transaction.StagedTable thread them through the rebuilt table, and Refresh re-hydrates them from the reloaded table. Labels() returns nil when the catalog returned none, aliases the table's labels so callers must treat the result as read-only, and labels are excluded from Table.Equals and View.Equals as transient state. Co-authored-by: Isaac --- catalog/rest/rest.go | 20 +-- catalog/rest/rest_test.go | 262 ++++++++++++++++++++++++++++++++++++ table/labels_wiring_test.go | 151 +++++++++++++++++++++ table/table.go | 23 ++++ table/transaction.go | 1 + view/labels_wiring_test.go | 43 ++++++ view/view.go | 37 ++++- 7 files changed, 527 insertions(+), 10 deletions(-) create mode 100644 table/labels_wiring_test.go create mode 100644 view/labels_wiring_test.go 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(