zeroshade commented on code in PR #2075:
URL: https://github.com/apache/iceberg-go/pull/2075#discussion_r4168212762


##########
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()
+       for k, v := range savedConfig {
+               assert.Equal(t, v, refreshedConfig[k], "saved config key %q", k)
+       }
+       assert.Empty(t, refreshedConfig["s3.access-key-id"])

Review Comment:
   This would still pass if the refresh saved `r.props` (which includes the 
catalog `token`) together with the caller's config. Comparing the whole map 
would catch that, and makes the loop above redundant:
   
   ```suggestion
        assert.Equal(t, savedConfig, refreshedConfig,
                "saved config must round-trip unchanged: no catalog props, no 
vended credentials")
   ```
   
   The same applies in `TestRefreshTableCredentialsAfterLoadTable`. After the 
load, add `assert.NotContains(t, tbl.SavedConfig(), "token")`, and after the 
refresh, add `assert.Equal(t, tbl.SavedConfig(), refreshed.SavedConfig())`.



##########
catalog/rest/rest.go:
##########
@@ -1303,6 +1309,41 @@ 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.
+//
+// Requires that the passed-in table instance be created with the 
table.WithSavedConfig() option to
+// save any table-specific configs. All tables created by this catalog pass in 
that option.

Review Comment:
   This isn't true for `UpdateTable`: it calls `tableFromResponse` without 
`WithSavedConfig` (`rest.go:1885-1889`). A refresh of that table therefore 
drops the metadata properties. Passing 
`table.WithSavedConfig(metadata.Properties())` there would fix it. Please also 
spell out here that `table.New` callers should include the metadata properties.



##########
table/table.go:
##########
@@ -1423,6 +1427,21 @@ func WithLabels(l *iceberg.Labels) Option {
        }
 }
 
+// WithSavedConfig supplies a set of properties used to create a *Table

Review Comment:
   Nit: please document the contract from the thread here and on 
`SavedConfig()`. The map holds table-scoped FileIO config only, never catalog 
props or credentials. It is visible to anyone who holds the table, and the REST 
refresh layers it between the catalog props and the vended creds.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to