Skip to content
Open
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
20 changes: 12 additions & 8 deletions catalog/rest/rest.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -1302,6 +1303,7 @@ func (r *Catalog) tableFromResponse(
r,
table.WithMetricsReporter(reporter),
table.WithScanPlanningIOProperties(scanPlanningConfig),
table.WithLabels(labels),
), nil
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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].
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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:"-"`
}

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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) {
Expand Down
195 changes: 195 additions & 0 deletions catalog/rest/rest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1633,6 +1633,56 @@ 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) TestLoadTableWithSnapshotModeRefs() {
Expand Down Expand Up @@ -2041,6 +2091,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)
Expand Down Expand Up @@ -2568,6 +2658,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() {
Expand Down Expand Up @@ -2953,6 +3092,62 @@ 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"
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"}}]}}`,
"metadata-location", 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) TestRegisterView404() {
const (
ns = "nonexistent"
Expand Down
20 changes: 20 additions & 0 deletions table/table.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -134,6 +137,10 @@ 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.
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 {
Expand Down Expand Up @@ -239,6 +246,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
Expand Down Expand Up @@ -1400,6 +1408,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
Expand Down
Loading
Loading