This is an automated email from the ASF dual-hosted git repository.
laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-go.git
The following commit(s) were added to refs/heads/main by this push:
new 8ff80185a feat(rest): add RefreshTableCredentials call to rest catalog
(#2075)
8ff80185a is described below
commit 8ff80185af86555cf99e5844a62cf6e31c271efd
Author: Bilal Akhtar <[email protected]>
AuthorDate: Mon Oct 5 15:03:42 2026 -0400
feat(rest): add RefreshTableCredentials call to rest catalog (#2075)
* feat(catalog): add RefreshTableCredentials call to catalog
This change exposes the table credential REST catalog endpoint as a
top-level method on catalogs. This is useful for when a table.Table is
seeded with a specific metadata location directly using table.New, but
is lacking temporarily-vended storage credentials to do any file scans.
* CR suggestion + one more rebase fix
* CR suggestions
* nit on comment
* another nit
---
catalog/catalog.go | 32 ++
catalog/rest/refresh_table_credentials_test.go | 468 +++++++++++++++++++++++++
catalog/rest/rest.go | 69 +++-
table/metrics_wiring_test.go | 25 ++
table/table.go | 19 +
table/transaction.go | 1 +
6 files changed, 605 insertions(+), 9 deletions(-)
diff --git a/catalog/catalog.go b/catalog/catalog.go
index a075b6cd3..1a508b320 100644
--- a/catalog/catalog.go
+++ b/catalog/catalog.go
@@ -213,6 +213,38 @@ type Closer interface {
Close() error
}
+// RefreshableCredentialCatalog is an optional interface implemented by
catalogs that
+// support temporarily-vended credentials. Temporarily-vended credentials
scope object
+// storage access to just files allowed by the user making a related request
to the
+// catalog. Callers should check for this capability via a type assertion:
+//
+// if rc, ok := cat.(catalog.RefreshableCredentialCatalog); ok {
+// refreshed, err := rc.RefreshTableCredentials(ctx, tbl)
+// if err == nil {
+// tbl = refreshed
+// }
+// }
+//
+// Passing the type assertion does not guarantee that credentials can be
vended.
+// An implementation may still return an error, for example when the catalog
+// service does not support vending credentials (the REST catalog returns
+// rest.ErrEndpointNotSupported if the server does not advertise the
credentials
+// endpoint) or no longer knows the table ([ErrNoSuchTable]). On error the
+// returned table is nil, and callers should keep using the original table.
+type RefreshableCredentialCatalog interface {
+ // RefreshTableCredentials loads temporarily-vended credentials for a
previously
+ // loaded table, and returns a new/modified table instance with
up-to-date credentials
+ // if the catalog is using temporarily-vended credentials for this
table.
+ // In this case, the instance should be able to refresh its own
credentials internally
+ // making it useful for use as a long-lived instance. For tables that do
+ // not have catalog-vended storage credentials, the table can be
returned with no
+ // modifications.
+ //
+ // The main use-case for this method is to be able to make
externally-created table
+ // instances compatible with catalog-based credential refreshes.
+ RefreshTableCredentials(ctx context.Context, tbl *table.Table)
(*table.Table, error)
+}
+
func ToIdentifier(ident ...string) table.Identifier {
if len(ident) == 1 {
if ident[0] == "" {
diff --git a/catalog/rest/refresh_table_credentials_test.go
b/catalog/rest/refresh_table_credentials_test.go
new file mode 100644
index 000000000..81cb40e60
--- /dev/null
+++ b/catalog/rest/refresh_table_credentials_test.go
@@ -0,0 +1,468 @@
+// 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 rest_test
+
+import (
+ "context"
+ "encoding/json"
+ "maps"
+ "net/http"
+ "net/http/httptest"
+ "net/url"
+ "strings"
+ "sync"
+ "sync/atomic"
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ "github.com/apache/iceberg-go/catalog"
+ "github.com/apache/iceberg-go/catalog/rest"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/metrics"
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// ioPropsRecorder registers an IO scheme that records the properties handed to
+// every filesystem load. That is how these tests observe whether vended
+// credentials actually reached a table's FileIO, which is otherwise private to
+// the table.
+type ioPropsRecorder struct {
+ mu sync.Mutex
+ loads []map[string]string
+}
+
+func (rec *ioPropsRecorder) register(t *testing.T, scheme string) {
+ t.Helper()
+
+ iceio.Register(scheme, func(_ context.Context, _ *url.URL, props
map[string]string) (iceio.IO, error) {
+ rec.mu.Lock()
+ defer rec.mu.Unlock()
+ rec.loads = append(rec.loads, maps.Clone(props))
+
+ return iceio.NewMemFS(), nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+}
+
+// lastLoad returns the properties the most recent filesystem load was built
+// from, failing the test if no load happened.
+func (rec *ioPropsRecorder) lastLoad(t *testing.T) map[string]string {
+ t.Helper()
+
+ rec.mu.Lock()
+ defer rec.mu.Unlock()
+ require.NotEmpty(t, rec.loads, "no filesystem was loaded")
+
+ return rec.loads[len(rec.loads)-1]
+}
+
+type credsCatalogOpts struct {
+ // endpoints is what /v1/config advertises; nil advertises every
endpoint.
+ endpoints []string
+ // defaults is the /v1/config defaults block, which the catalog folds
into
+ // the properties it later builds table FileIO from.
+ defaults map[string]string
+ // creds handles GET /v1/namespaces/db/tables/tbl/credentials. When
nil, the
+ // endpoint is left unrouted so a request to it would 404.
+ creds http.HandlerFunc
+ // loadTable handles GET /v1/namespaces/db/tables/tbl. When nil, the
+ // endpoint is left unrouted so a request to it would 404.
+ loadTable http.HandlerFunc
+}
+
+func newCredsTestCatalog(t *testing.T, opts credsCatalogOpts) *rest.Catalog {
+ t.Helper()
+
+ mux := http.NewServeMux()
+ mux.HandleFunc("/v1/config", func(w http.ResponseWriter, _
*http.Request) {
+ endpoints := opts.endpoints
+ if endpoints == nil {
+ endpoints = rest.AllEndpointStrings
+ }
+ defaults := opts.defaults
+ if defaults == nil {
+ defaults = map[string]string{}
+ }
+ assert.NoError(t, json.NewEncoder(w).Encode(map[string]any{
+ "defaults": defaults,
+ "overrides": map[string]any{},
+ "endpoints": endpoints,
+ }))
+ })
+ if opts.creds != nil {
+ mux.HandleFunc("/v1/namespaces/db/tables/tbl/credentials",
opts.creds)
+ }
+ if opts.loadTable != nil {
+ mux.HandleFunc("/v1/namespaces/db/tables/tbl", opts.loadTable)
+ }
+
+ srv := httptest.NewServer(mux)
+ t.Cleanup(srv.Close)
+
+ cat, err := rest.NewCatalog(context.Background(), "rest", srv.URL,
rest.WithOAuthToken(TestToken))
+ require.NoError(t, err)
+ t.Cleanup(func() { assert.NoError(t, cat.Close()) })
+
+ return cat
+}
+
+// newExternalTable builds a table the way a caller outside the catalog would:
+// straight through table.New, with a plain property-based FileIO and no vended
+// credentials seeded into it. Any opts are passed through to table.New.
+func newExternalTable(t *testing.T, cat *rest.Catalog, scheme string, opts
...table.Option) (*table.Table, string) {
+ t.Helper()
+
+ // The shared fixture is written against s3://; point it at the
recording
+ // scheme so loading its FileIO stays offline and observable.
+ meta, err := table.ParseMetadataString(
+ strings.ReplaceAll(exampleTableMetadataNoSnapshotV1, "s3://",
scheme+"://"))
+ require.NoError(t, err)
+
+ metadataLoc := scheme +
"://warehouse/database/table/metadata/00000-a.metadata.json"
+
+ return table.New(
+ catalog.ToIdentifier("db", "tbl"),
+ meta,
+ metadataLoc,
+ iceio.LoadFSFunc(nil, metadataLoc),
+ cat,
+ opts...,
+ ), metadataLoc
+}
+
+func storageCredentialsBody(prefix string, config map[string]string)
map[string]any {
+ return map[string]any{
+ "storage-credentials": []any{
+ map[string]any{"prefix": prefix, "config": config},
+ },
+ }
+}
+
+// TestRefreshTableCredentialsSeedsExternallyCreatedTable covers the case a
+// caller cannot reach through LoadTable: a table handed to table.New directly,
+// whose FileIO was never seeded with vended credentials, is brought up to date
+// without a metadata reload.
+func TestRefreshTableCredentialsSeedsExternallyCreatedTable(t *testing.T) {
+ const scheme = "restcreds-seed"
+
+ rec := &ioPropsRecorder{}
+ rec.register(t, scheme)
+
+ var credsCalls atomic.Int32
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ defaults: map[string]string{"s3.region": "us-west-2"},
+ creds: func(w http.ResponseWriter, req *http.Request) {
+ credsCalls.Add(1)
+ assert.Equal(t, http.MethodGet, req.Method)
+ assert.NoError(t,
json.NewEncoder(w).Encode(storageCredentialsBody(
+ scheme+"://warehouse/database/table",
+ map[string]string{
+ "s3.access-key-id": "vended-key",
+ "s3.secret-access-key": "vended-secret",
+ "s3.session-token": "vended-token",
+ },
+ )))
+ },
+ })
+
+ // Save config the way the catalog does for tables it loads: the table's
+ // metadata properties (the fixture sets this codec) plus table-specific
+ // FileIO settings such as a region or client factory.
+ savedConfig := iceberg.Properties{
+ "write.parquet.compression-codec": "zstd",
+ // Overrides the catalog default above.
+ iceio.S3Region: "eu-central-1",
+ "client.factory": "com.example.CustomClientFactory",
+ }
+ tbl, metadataLoc := newExternalTable(t, cat, scheme,
table.WithSavedConfig(savedConfig))
+
+ // The premise: as built, the table's FileIO carries no credentials.
+ _, err := tbl.FS(context.Background())
+ require.NoError(t, err)
+ assert.NotContains(t, rec.lastLoad(t), "s3.access-key-id")
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ require.NotNil(t, refreshed)
+ assert.Equal(t, int32(1), credsCalls.Load())
+
+ // Only the FileIO configuration changes: identity, metadata and the
metadata
+ // location a later commit targets all carry over untouched.
+ assert.Equal(t, catalog.ToIdentifier("db", "tbl"),
refreshed.Identifier())
+ assert.Equal(t, metadataLoc, refreshed.MetadataLocation())
+ assert.Equal(t, tbl.Metadata(), refreshed.Metadata())
+ assert.Equal(t, tbl.Location(), refreshed.Location())
+
+ _, err = refreshed.FS(context.Background())
+ require.NoError(t, err)
+
+ props := rec.lastLoad(t)
+ assert.Equal(t, "vended-key", props["s3.access-key-id"])
+ assert.Equal(t, "vended-secret", props["s3.secret-access-key"])
+ assert.Equal(t, "vended-token", props["s3.session-token"])
+ // The table's saved config survives the merge, winning over the
catalog's
+ // own config, rather than being replaced by the credentials alone.
+ for k, v := range savedConfig {
+ assert.Equal(t, v, props[k], "FileIO property %q", k)
+ }
+
+ // The saved config also carries over to the refreshed table itself,
without
+ // any changes - the vended credentials are not merged in it.
+ refreshedConfig := refreshed.SavedConfig()
+ assert.Equal(t, savedConfig, refreshedConfig,
+ "saved config must round-trip unchanged: no catalog props, no
vended credentials")
+
+ // The table the caller passed in, and the map it saved, are left alone.
+ _, err = tbl.FS(context.Background())
+ require.NoError(t, err)
+ assert.NotContains(t, rec.lastLoad(t), "s3.access-key-id")
+ assert.Equal(t, savedConfig, tbl.SavedConfig())
+ assert.NotContains(t, savedConfig, "s3.access-key-id")
+}
+
+// customReporter is a non-nop metrics reporter a caller might attach to a
+// table, distinguishable from the catalog's default by pointer identity.
+type customReporter struct{}
+
+func (*customReporter) Report(context.Context, metrics.MetricsReport) {}
+func (*customReporter) Close() error { return
nil }
+
+// TestRefreshTableCredentialsPreservesMetricsReporter checks that a reporter
+// the caller set on the table survives the refresh instead of being replaced
+// by the catalog's default reporter.
+func TestRefreshTableCredentialsPreservesMetricsReporter(t *testing.T) {
+ const scheme = "restcreds-reporter"
+
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ assert.NoError(t,
json.NewEncoder(w).Encode(storageCredentialsBody(
+ scheme+"://warehouse/database/table",
+ map[string]string{"s3.access-key-id":
"vended-key"},
+ )))
+ },
+ })
+
+ reporter := &customReporter{}
+ tbl, _ := newExternalTable(t, cat, scheme,
table.WithMetricsReporter(reporter))
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ require.NotSame(t, tbl, refreshed)
+ assert.Same(t, reporter, refreshed.MetricsReporter())
+}
+
+// TestRefreshTableCredentialsAfterLoadTable covers a table loaded through
+// LoadTable: the per-table config block of the load response is in neither the
+// catalog's props nor the table's metadata properties, so only the table's
+// saved config can carry it through a credential refresh.
+func TestRefreshTableCredentialsAfterLoadTable(t *testing.T) {
+ const scheme = "restcreds-loaded"
+
+ rec := &ioPropsRecorder{}
+ rec.register(t, scheme)
+
+ metadataLoc := scheme +
"://warehouse/database/table/metadata/00000-a.metadata.json"
+ tableConfig := map[string]string{
+ // Overrides the catalog default below.
+ iceio.S3Region: "eu-central-1",
+ "client.factory": "com.example.CustomClientFactory",
+ }
+
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ defaults: map[string]string{iceio.S3Region: "us-west-2"},
+ loadTable: func(w http.ResponseWriter, req *http.Request) {
+ assert.Equal(t, http.MethodGet, req.Method)
+ assert.NoError(t,
json.NewEncoder(w).Encode(map[string]any{
+ "metadata-location": metadataLoc,
+ "metadata": json.RawMessage(strings.ReplaceAll(
+ exampleTableMetadataNoSnapshotV1,
"s3://", scheme+"://")),
+ "config": tableConfig,
+ }))
+ },
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ assert.NoError(t,
json.NewEncoder(w).Encode(storageCredentialsBody(
+ scheme+"://warehouse/database/table",
+ map[string]string{
+ iceio.S3AccessKeyID: "vended-key",
+ iceio.S3SecretAccessKey:
"vended-secret",
+ iceio.S3SessionToken: "vended-token",
+ },
+ )))
+ },
+ })
+
+ tbl, err := cat.LoadTable(context.Background(),
catalog.ToIdentifier("db", "tbl"))
+ require.NoError(t, err)
+
+ // The premise: the loaded table's FileIO carries the per-table config
but no
+ // credentials, since the load response vended none.
+ _, err = tbl.FS(context.Background())
+ require.NoError(t, err)
+ props := rec.lastLoad(t)
+ for k, v := range tableConfig {
+ require.Equal(t, v, props[k], "loaded FileIO property %q", k)
+ }
+ require.NotContains(t, props, iceio.S3AccessKeyID)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ require.NotNil(t, refreshed)
+ assert.Equal(t, tbl.Identifier(), refreshed.Identifier())
+ assert.Equal(t, metadataLoc, refreshed.MetadataLocation())
+ assert.Equal(t, tbl.Metadata(), refreshed.Metadata())
+
+ _, err = refreshed.FS(context.Background())
+ require.NoError(t, err)
+
+ props = rec.lastLoad(t)
+ for k, v := range tableConfig {
+ assert.Equal(t, v, props[k], "refreshed FileIO property %q", k)
+ }
+ assert.Equal(t, "zstd", props["write.parquet.compression-codec"],
+ "metadata properties merged at load time should survive too")
+ assert.Equal(t, "vended-key", props[iceio.S3AccessKeyID])
+ assert.Equal(t, "vended-secret", props[iceio.S3SecretAccessKey])
+ assert.Equal(t, "vended-token", props[iceio.S3SessionToken])
+ // Saved configs are identical.
+ assert.Equal(t, tbl.SavedConfig(), refreshed.SavedConfig())
+}
+
+// TestRefreshTableCredentialsSeededTableRenewsOnExpiry checks the refreshed
+// table's FileIO keeps the refresher wiring, so credentials that expire are
+// re-fetched instead of being pinned at their first value.
+func TestRefreshTableCredentialsSeededTableRenewsOnExpiry(t *testing.T) {
+ const scheme = "restcreds-renew"
+
+ rec := &ioPropsRecorder{}
+ rec.register(t, scheme)
+
+ var credsCalls atomic.Int32
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ n := credsCalls.Add(1)
+ assert.NoError(t,
json.NewEncoder(w).Encode(storageCredentialsBody(
+ scheme+"://warehouse/database/table",
+ map[string]string{
+ "s3.access-key-id": "vended-key",
+ // Already past its expiry, so the very
next FileIO load must
+ // go back to the catalog rather than
reuse it.
+ "s3.session-token-expires-at-ms": "1",
+ "s3.session-token":
"token-" + strings.Repeat("x", int(n)),
+ },
+ )))
+ },
+ })
+
+ tbl, _ := newExternalTable(t, cat, scheme)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ require.Equal(t, int32(1), credsCalls.Load())
+
+ _, err = refreshed.FS(context.Background())
+ require.NoError(t, err)
+ assert.Equal(t, "token-x", rec.lastLoad(t)["s3.session-token"])
+
+ _, err = refreshed.FS(context.Background())
+ require.NoError(t, err)
+ assert.Equal(t, int32(2), credsCalls.Load(),
+ "expired credentials should be renewed through the catalog")
+ assert.Equal(t, "token-xx", rec.lastLoad(t)["s3.session-token"])
+}
+
+func TestRefreshTableCredentialsNoCredentialsVended(t *testing.T) {
+ const scheme = "restcreds-empty"
+
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ assert.NoError(t,
json.NewEncoder(w).Encode(map[string]any{
+ "storage-credentials": []any{},
+ }))
+ },
+ })
+
+ tbl, _ := newExternalTable(t, cat, scheme)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ assert.Same(t, tbl, refreshed,
+ "a server that vends nothing should leave the table as-is")
+}
+
+// TestRefreshTableCredentialsPrefixMismatch pins that credentials are resolved
+// against the table's metadata location: a credential for some other prefix
+// must not be applied.
+func TestRefreshTableCredentialsPrefixMismatch(t *testing.T) {
+ const scheme = "restcreds-mismatch"
+
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ assert.NoError(t,
json.NewEncoder(w).Encode(storageCredentialsBody(
+
scheme+"://warehouse/other-database/other-table",
+ map[string]string{"s3.access-key-id":
"vended-key"},
+ )))
+ },
+ })
+
+ tbl, _ := newExternalTable(t, cat, scheme)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ require.NoError(t, err)
+ assert.Same(t, tbl, refreshed,
+ "a server that vends a different prefix credential should leave
the table as-is")
+}
+
+func TestRefreshTableCredentialsEndpointNotAdvertised(t *testing.T) {
+ const scheme = "restcreds-unsupported"
+
+ // Advertise load-table but not the credentials endpoint.
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ endpoints: []string{"GET
/v1/{prefix}/namespaces/{namespace}/tables/{table}"},
+ })
+
+ tbl, _ := newExternalTable(t, cat, scheme)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ assert.Nil(t, refreshed)
+ assert.ErrorIs(t, err, rest.ErrEndpointNotSupported)
+}
+
+func TestRefreshTableCredentialsNotFound(t *testing.T) {
+ const scheme = "restcreds-notfound"
+
+ cat := newCredsTestCatalog(t, credsCatalogOpts{
+ creds: func(w http.ResponseWriter, _ *http.Request) {
+ w.WriteHeader(http.StatusNotFound)
+ assert.NoError(t,
json.NewEncoder(w).Encode(map[string]any{
+ "error": map[string]any{
+ "message": "Table does not exist:
db.tbl",
+ "type": "NoSuchTableException",
+ "code": 404,
+ },
+ }))
+ },
+ })
+
+ tbl, _ := newExternalTable(t, cat, scheme)
+
+ refreshed, err := cat.RefreshTableCredentials(context.Background(), tbl)
+ assert.Nil(t, refreshed)
+ assert.ErrorIs(t, err, catalog.ErrNoSuchTable)
+}
diff --git a/catalog/rest/rest.go b/catalog/rest/rest.go
index fae679108..2ad477634 100644
--- a/catalog/rest/rest.go
+++ b/catalog/rest/rest.go
@@ -49,8 +49,9 @@ import (
)
var (
- _ catalog.Catalog = (*Catalog)(nil)
- _ catalog.TransactionalCatalog = (*Catalog)(nil)
+ _ catalog.Catalog = (*Catalog)(nil)
+ _ catalog.TransactionalCatalog = (*Catalog)(nil)
+ _ catalog.RefreshableCredentialCatalog = (*Catalog)(nil)
)
const (
@@ -1217,6 +1218,7 @@ func (r *Catalog) tableFromResponse(
scanPlanningConfig iceberg.Properties,
credsVended bool,
labels *iceberg.Labels,
+ opts ...table.Option,
) (*table.Table, error) {
var fsF func(context.Context) (iceio.IO, error)
if credsVended {
@@ -1267,15 +1269,19 @@ func (r *Catalog) tableFromResponse(
}
}
+ // Caller-supplied opts are applied last so they can override the
defaults
+ // derived above.
return table.New(
identifier,
metadata,
loc,
fsF,
r,
- table.WithMetricsReporter(reporter),
- table.WithScanPlanningIOProperties(scanPlanningConfig),
- table.WithLabels(labels),
+ append([]table.Option{
+ table.WithMetricsReporter(reporter),
+ table.WithScanPlanningIOProperties(scanPlanningConfig),
+ table.WithLabels(labels),
+ }, opts...)...,
), nil
}
@@ -1303,6 +1309,43 @@ func (r *Catalog) fetchTableCreds(ctx context.Context,
ident []string, location
return resolveStorageCredentials(ret.StorageCredentials, location), nil
}
+// RefreshTableCredentials updates a *table.Table with newly-vended
credentials from the catalog
+// without updating any other table-internal state that a full Refresh() would.
+// Allows for a quick table credential refresh if the table was created
without any pre-seeded
+// credentials. If the catalog did not vend any credentials, the table is
returned unmodified.
+//
+// Note that the passed-in table instance should be created with the
table.WithSavedConfig() option to
+// save any table-specific configs for reuse in the newly-created instance.
All tables created by this
+// catalog pass in that option. Callers of table.New() that use this method
should also pass
+// metadata.Properties() into table.WithSavedConfig().
+func (r *Catalog) RefreshTableCredentials(ctx context.Context, tbl
*table.Table) (*table.Table, error) {
+ metadataLoc := tbl.MetadataLocation()
+ resp, err := r.fetchTableCreds(ctx, tbl.Identifier(), metadataLoc)
+ if err != nil {
+ return nil, err
+ }
+ if len(resp) == 0 {
+ // No new credentials vended. Return as-is.
+ return tbl, nil
+ }
+
+ // Return a new *table.Table with newly-merged credentials coming from
the
+ // fetchTableCreds call.
+ config := maps.Clone(r.props)
+ maps.Copy(config, tbl.SavedConfig())
+ scanCfg := maps.Clone(config)
+ maps.Copy(config, resp)
+
+ // Keep a reporter the caller set on the table rather than reverting it
to
+ // the catalog default, as Refresh does.
+ opts := []table.Option{table.WithSavedConfig(tbl.SavedConfig())}
+ if reporter := tbl.MetricsReporter(); !metrics.IsNop(reporter) {
+ opts = append(opts, table.WithMetricsReporter(reporter))
+ }
+
+ return r.tableFromResponse(ctx, tbl.Identifier(), tbl.Metadata(),
metadataLoc, config, scanCfg, true, tbl.Labels(), opts...)
+}
+
type identifierPageFetcher func(pageToken string) ([]table.Identifier, string,
error)
func paginateIdentifiers(fetch identifierPageFetcher, yield
func(table.Identifier) bool) error {
@@ -1517,12 +1560,15 @@ func (r *Catalog) CreateTable(ctx context.Context,
identifier table.Identifier,
config := maps.Clone(r.props)
maps.Copy(config, ret.Metadata.Properties())
+ // Save only the per-table configs.
+ saved := maps.Clone(ret.Metadata.Properties())
maps.Copy(config, ret.Config)
+ maps.Copy(saved, ret.Config)
scanPlanningConfig := maps.Clone(config)
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, ret.Labels)
+ return r.tableFromResponse(ctx, identifier, ret.Metadata,
ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels,
table.WithSavedConfig(saved))
}
// commitStagedCreate performs the second phase of a staged table
@@ -1745,12 +1791,14 @@ func (r *Catalog) RegisterTable(ctx context.Context,
identifier table.Identifier
config := maps.Clone(r.props)
maps.Copy(config, ret.Metadata.Properties())
+ saved := maps.Clone(ret.Metadata.Properties())
maps.Copy(config, ret.Config)
+ maps.Copy(saved, ret.Config)
scanPlanningConfig := maps.Clone(config)
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, ret.Labels)
+ return r.tableFromResponse(ctx, identifier, ret.Metadata,
ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels,
table.WithSavedConfig(saved))
}
// LoadTable loads a table from the catalog. It implements [catalog.Catalog].
@@ -1790,12 +1838,15 @@ func (r *Catalog) loadTableWithMode(ctx
context.Context, identifier table.Identi
config := maps.Clone(r.props)
maps.Copy(config, ret.Metadata.Properties())
+ // Save only the per-table configs, not r.props.
+ saved := maps.Clone(ret.Metadata.Properties())
maps.Copy(config, ret.Config)
+ maps.Copy(saved, ret.Config)
scanPlanningConfig := maps.Clone(config)
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, ret.Labels)
+ return r.tableFromResponse(ctx, identifier, ret.Metadata,
ret.MetadataLoc, config, scanPlanningConfig, credsVended, ret.Labels,
table.WithSavedConfig(saved))
}
func (r *Catalog) UpdateTable(ctx context.Context, ident table.Identifier,
requirements []table.Requirement, updates []table.Update) (*table.Table, error)
{
@@ -1844,7 +1895,7 @@ func (r *Catalog) UpdateTable(ctx context.Context, ident
table.Identifier, requi
maps.Copy(config, metadata.Properties())
// A commit response carries no labels (they are load-time enrichment).
- return r.tableFromResponse(ctx, ident, metadata, ret.MetadataLoc,
config, config, false, nil)
+ return r.tableFromResponse(ctx, ident, metadata, ret.MetadataLoc,
config, config, false, nil, table.WithSavedConfig(metadata.Properties()))
}
func (r *Catalog) DropTable(ctx context.Context, identifier table.Identifier)
error {
diff --git a/table/metrics_wiring_test.go b/table/metrics_wiring_test.go
index 8f1cec26a..86633dc1d 100644
--- a/table/metrics_wiring_test.go
+++ b/table/metrics_wiring_test.go
@@ -116,6 +116,31 @@ func TestStagedTableInheritsReporter(t *testing.T) {
"StagedTable must forward the transaction table's reporter")
}
+// TestStagedTableInheritsSavedConfig pins that StagedTable forwards the saved
+// config, so a credential refresh of a staged table, or of anything committed
+// from it, still sees the config the original table was created with.
+func TestStagedTableInheritsSavedConfig(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)
+
+ savedConfig := iceberg.Properties{
+ "s3.region": "eu-central-1",
+ "client.factory": "com.example.CustomClientFactory",
+ }
+ tbl := New(Identifier{"default", "staged"}, meta, "",
+ func(context.Context) (iceio.IO, error) { return
iceio.NewMemFS(), nil }, nil,
+ WithSavedConfig(savedConfig))
+
+ staged, err := tbl.NewTransaction().StagedTable()
+ require.NoError(t, err)
+ assert.Equal(t, savedConfig, staged.SavedConfig(),
+ "StagedTable must forward the transaction table's saved config")
+}
+
// TestRefreshKeepsCallerReporter pins that a reporter injected via
// WithMetricsReporter must survive a Refresh even though the catalog's
LoadTable
// hands back the default nop reporter.
diff --git a/table/table.go b/table/table.go
index 6771c6ad8..c77d66517 100644
--- a/table/table.go
+++ b/table/table.go
@@ -111,6 +111,7 @@ type Table struct {
// REST catalogs can return FileIO configuration in the response's
config
// block.
scanPlanningIOProps iceberg.Properties
+ savedConfig iceberg.Properties
reporter metrics.Reporter
// reporterSet records whether a caller injected a reporter via
// WithMetricsReporter. It distinguishes an explicit reporter
(including an
@@ -136,6 +137,7 @@ func (t Table) Schema() *iceberg.Schema
{ return t.metadata
func (t Table) Spec() iceberg.PartitionSpec { return
t.metadata.PartitionSpec() }
func (t Table) SortOrder() SortOrder { return
t.metadata.SortOrder() }
func (t Table) Properties() iceberg.Properties { return
t.metadata.Properties() }
+func (t Table) SavedConfig() iceberg.Properties { return
maps.Clone(t.savedConfig) }
// Labels returns the catalog-provided labels from the load response, or nil if
// the catalog returned none. Labels are transient enrichment, not table state.
@@ -249,6 +251,7 @@ func (t *Table) Refresh(ctx context.Context) error {
t.planner = fresh.planner
t.scanPlanningIOProps = maps.Clone(fresh.scanPlanningIOProps)
t.labels = fresh.labels
+ t.savedConfig = maps.Clone(fresh.savedConfig)
// 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
@@ -836,6 +839,7 @@ func (t Table) doCommit(ctx context.Context, updates
[]Update, reqs []Requiremen
withReporterState(t.reporter, t.reporterSet),
WithScanPlanningIOProperties(t.scanPlanningIOProps),
WithLabels(t.labels),
+ WithSavedConfig(t.savedConfig),
), nil
}
@@ -1423,6 +1427,21 @@ func WithLabels(l *iceberg.Labels) Option {
}
}
+// WithSavedConfig supplies a set of properties used to create a *Table
+// instance. This saved config is exposed through table.SavedConfig() and
+// can be used by callers to save properties along with a table instance, for
+// reuse later if the table instance were to be cloned in a new call to
+// table.New.
+func WithSavedConfig(config iceberg.Properties) Option {
+ if config == nil {
+ return noopTableOption
+ }
+
+ return func(t *Table) {
+ t.savedConfig = maps.Clone(config)
+ }
+}
+
// 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 b83acc33d..94817657b 100644
--- a/table/transaction.go
+++ b/table/transaction.go
@@ -3328,6 +3328,7 @@ func (t *Transaction) StagedTable() (*StagedTable, error)
{
withReporterState(t.tbl.reporter, t.tbl.reporterSet),
WithScanPlanningIOProperties(t.tbl.scanPlanningIOProps),
WithLabels(t.tbl.labels),
+ WithSavedConfig(t.tbl.savedConfig),
),
}, nil
}