-
Notifications
You must be signed in to change notification settings - Fork 249
fix(catalog/rest): sign SigV4 with credentials from catalog properties #1999
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
ffa55bb
cca7c83
48164f5
daa7aa3
1b37079
d401049
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -88,6 +88,14 @@ func WithMetadataLocation(loc string) Option { | |
| } | ||
| } | ||
|
|
||
| // WithSigV4 enables AWS SigV4 request signing for the REST catalog. The signing | ||
| // identity is resolved in order: an explicit WithAwsConfig, then the s3.* catalog | ||
| // credential properties (s3.access-key-id / s3.secret-access-key / s3.session-token), | ||
| // then the AWS default credential chain. | ||
| // | ||
| // The Java-client property names (rest.access-key-id / rest.secret-access-key / | ||
| // rest.session-token) are accepted as aliases, resolved per field with the s3.* | ||
| // keys taking precedence when both are set. | ||
|
Comment on lines
+96
to
+98
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This is stale after d401049: resolution is no longer per field. Suggested text: "If any s3.* credential key is set, only the s3.* tuple is used and it must be complete; otherwise the rest.* tuple is used." |
||
| func WithSigV4() Option { | ||
| return func(o *options) { | ||
| o.enableSigv4 = true | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -41,6 +41,7 @@ import ( | |
|
|
||
| "github.com/apache/iceberg-go" | ||
| "github.com/apache/iceberg-go/catalog" | ||
| internalaws "github.com/apache/iceberg-go/internal/awsconfig" | ||
| iceio "github.com/apache/iceberg-go/io" | ||
| "github.com/apache/iceberg-go/metrics" | ||
| "github.com/apache/iceberg-go/table" | ||
|
|
@@ -49,6 +50,7 @@ import ( | |
| "github.com/aws/aws-sdk-go-v2/aws" | ||
| v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4" | ||
| "github.com/aws/aws-sdk-go-v2/config" | ||
| "github.com/aws/aws-sdk-go-v2/credentials" | ||
| "golang.org/x/oauth2" | ||
| "golang.org/x/oauth2/clientcredentials" | ||
| "golang.org/x/sync/semaphore" | ||
|
|
@@ -92,6 +94,12 @@ const ( | |
| keyRestSigV4Region = "rest.signing-region" | ||
| keyRestSigV4Service = "rest.signing-name" | ||
| keyAuthUrl = "rest.authorization-url" | ||
| // keyRestAccessKeyID and friends are the Java-client property names for the | ||
| // SigV4 signing credentials. They are accepted as aliases for the s3.* | ||
| // properties; the s3.* keys take precedence when both are set. | ||
| keyRestAccessKeyID = "rest.access-key-id" | ||
| keyRestSecretAccessKey = "rest.secret-access-key" | ||
| keyRestSessionToken = "rest.session-token" | ||
| // keyOAuth2ServerURI is the portable, spec-aligned property for the OAuth2 | ||
| // token endpoint used by Java, PyIceberg and iceberg-rust. It is the | ||
| // preferred key; keyAuthUrl is retained as a compatibility alias. When both | ||
|
|
@@ -242,6 +250,31 @@ type sessionTransport struct { | |
| cfg aws.Config | ||
| service string | ||
| newHash func() hash.Hash | ||
| // signingOrigin is the configured catalog origin. Requests to a different | ||
| // origin (e.g. a redirect hop) are not signed, so the SigV4 Authorization | ||
| // header and session token never reach an unconfigured host. | ||
| signingOrigin *url.URL | ||
| } | ||
|
|
||
| // sameOrigin reports whether two URLs share scheme, host, and effective port. | ||
| func sameOrigin(a, b *url.URL) bool { | ||
| return strings.EqualFold(a.Scheme, b.Scheme) && | ||
| strings.EqualFold(a.Hostname(), b.Hostname()) && | ||
| defaultedPort(a) == defaultedPort(b) | ||
| } | ||
|
|
||
| func defaultedPort(u *url.URL) string { | ||
| if p := u.Port(); p != "" { | ||
| return p | ||
| } | ||
| switch strings.ToLower(u.Scheme) { | ||
| case "https": | ||
| return "443" | ||
| case "http": | ||
| return "80" | ||
| default: | ||
| return "" | ||
| } | ||
| } | ||
|
|
||
| // from https://pkg.go.dev/github.com/aws/aws-sdk-go-v2/aws/signer/v4#Signer.SignHTTP | ||
|
|
@@ -291,7 +324,7 @@ func (s *sessionTransport) RoundTrip(r *http.Request) (*http.Response, error) { | |
| r.Header.Set(k, v) | ||
| } | ||
|
|
||
| if s.signer != nil { | ||
| if s.signer != nil && (s.signingOrigin == nil || sameOrigin(s.signingOrigin, r.URL)) { | ||
| var payloadHash string | ||
| if r.Body == nil { | ||
| payloadHash = emptyStringHash | ||
|
|
@@ -1106,26 +1139,64 @@ func (r *Catalog) createSession(ctx context.Context, opts *options) (*http.Clien | |
| if opts.enableSigv4 { | ||
| cfg := opts.awsConfig | ||
| if !opts.awsConfigSet { | ||
| creds, err := staticCredsFromProps(opts.additionalProps) | ||
|
iremcaginyurtturk marked this conversation as resolved.
|
||
| if err != nil { | ||
| cleanup() | ||
|
|
||
| return nil, nil, err | ||
| } | ||
| // If no config provided, load defaults from environment. | ||
| var err error | ||
| cfg, err = config.LoadDefaultConfig(ctx) | ||
| if err != nil { | ||
| cleanup() | ||
|
|
||
| return nil, nil, err | ||
| } | ||
| // Sign with the S3 credentials carried in the catalog properties when | ||
| // present, rather than only the AWS default credential chain. | ||
|
iremcaginyurtturk marked this conversation as resolved.
|
||
| if creds != nil { | ||
| cfg.Credentials = creds | ||
|
iremcaginyurtturk marked this conversation as resolved.
iremcaginyurtturk marked this conversation as resolved.
|
||
| } | ||
| } | ||
| if opts.sigv4Region != "" { | ||
| cfg.Region = opts.sigv4Region | ||
| } | ||
|
|
||
| session.cfg, session.service = cfg, opts.sigv4Service | ||
| session.signer, session.newHash = v4.NewSigner(), sha256.New | ||
| session.signingOrigin = r.baseURI | ||
| } | ||
|
|
||
| return cl, cleanup, nil | ||
| } | ||
|
|
||
| // staticCredsFromProps returns a static credentials provider built from the | ||
| // signing-credential properties. It prefers the s3.* keys and falls back to the | ||
| // Java-compatible rest.* aliases, resolving the tuple atomically from a single | ||
| // namespace so a partial pair is never completed with fields from the other one. | ||
| // It returns (nil, nil) when neither namespace sets any credential property, so | ||
| // the caller falls back to the default credential chain, and an | ||
| // ErrIncompleteStaticCredentials error when the chosen namespace is incomplete. | ||
| func staticCredsFromProps(props iceberg.Properties) (aws.CredentialsProvider, error) { | ||
|
iremcaginyurtturk marked this conversation as resolved.
|
||
| namespaces := [][3]string{ | ||
| {iceio.S3AccessKeyID, iceio.S3SecretAccessKey, iceio.S3SessionToken}, | ||
| {keyRestAccessKeyID, keyRestSecretAccessKey, keyRestSessionToken}, | ||
| } | ||
|
Comment on lines
+1181
to
+1184
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Non-blocking, but worth settling before this ships because it is user-visible. With both namespaces set, Java's I'd flip this to |
||
| for _, ns := range namespaces { | ||
| accessKey, secretKey, token := props[ns[0]], props[ns[1]], props[ns[2]] | ||
| if accessKey == "" && secretKey == "" && token == "" { | ||
| continue | ||
| } | ||
| if err := internalaws.ValidateStaticCredentials(ns[0], ns[1], ns[2], accessKey, secretKey, token); err != nil { | ||
| return nil, err | ||
| } | ||
|
|
||
| return credentials.NewStaticCredentialsProvider(accessKey, secretKey, token), nil | ||
| } | ||
|
|
||
| return nil, nil | ||
| } | ||
|
|
||
| func (r *Catalog) fetchConfig(ctx context.Context, opts *options) (*options, error) { | ||
| params := url.Values{} | ||
| if opts.warehouseLocation != "" { | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -42,6 +42,8 @@ import ( | |
|
|
||
| "github.com/apache/iceberg-go" | ||
| "github.com/apache/iceberg-go/catalog" | ||
| internalaws "github.com/apache/iceberg-go/internal/awsconfig" | ||
| iceio "github.com/apache/iceberg-go/io" | ||
| "github.com/apache/iceberg-go/table" | ||
| "github.com/aws/aws-sdk-go-v2/aws" | ||
| v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4" | ||
|
|
@@ -52,6 +54,159 @@ import ( | |
| "golang.org/x/sync/errgroup" | ||
| ) | ||
|
|
||
| func TestStaticCredsFromProps(t *testing.T) { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Optional, as @laskoviymishka noted earlier: this reads better as a table-driven test with |
||
| creds, err := staticCredsFromProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AK", | ||
| iceio.S3SecretAccessKey: "SK", | ||
| iceio.S3SessionToken: "ST", | ||
| }) | ||
| require.NoError(t, err) | ||
| require.NotNil(t, creds) | ||
| got, err := creds.Retrieve(context.Background()) | ||
| require.NoError(t, err) | ||
| require.Equal(t, "AK", got.AccessKeyID) | ||
| require.Equal(t, "SK", got.SecretAccessKey) | ||
| require.Equal(t, "ST", got.SessionToken) | ||
|
|
||
| creds, err = staticCredsFromProps(iceberg.Properties{}) | ||
| require.NoError(t, err, "no creds must fall back to the default chain") | ||
| require.Nil(t, creds) | ||
|
|
||
| _, err = staticCredsFromProps(iceberg.Properties{iceio.S3AccessKeyID: "AK"}) | ||
| require.ErrorIs(t, err, internalaws.ErrIncompleteStaticCredentials, "a lone access key must be an error, not the ambient identity") | ||
|
|
||
| _, err = staticCredsFromProps(iceberg.Properties{iceio.S3SecretAccessKey: "SK"}) | ||
|
iremcaginyurtturk marked this conversation as resolved.
|
||
| require.ErrorIs(t, err, internalaws.ErrIncompleteStaticCredentials, "a lone secret key must be an error, not the ambient identity") | ||
|
|
||
| _, err = staticCredsFromProps(iceberg.Properties{iceio.S3SessionToken: "ST"}) | ||
| require.ErrorIs(t, err, internalaws.ErrIncompleteStaticCredentials, "a lone session token must be an error, not the ambient identity") | ||
|
|
||
| creds, err = staticCredsFromProps(iceberg.Properties{ | ||
| keyRestAccessKeyID: "RAK", | ||
| keyRestSecretAccessKey: "RSK", | ||
| keyRestSessionToken: "RST", | ||
| }) | ||
| require.NoError(t, err) | ||
| require.NotNil(t, creds) | ||
| got, err = creds.Retrieve(context.Background()) | ||
| require.NoError(t, err) | ||
| require.Equal(t, "RAK", got.AccessKeyID) | ||
| require.Equal(t, "RSK", got.SecretAccessKey) | ||
| require.Equal(t, "RST", got.SessionToken) | ||
|
|
||
| creds, err = staticCredsFromProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AK", | ||
| iceio.S3SecretAccessKey: "SK", | ||
| keyRestAccessKeyID: "RAK", | ||
| keyRestSecretAccessKey: "RSK", | ||
| }) | ||
| require.NoError(t, err) | ||
| got, err = creds.Retrieve(context.Background()) | ||
| require.NoError(t, err) | ||
| require.Equal(t, "AK", got.AccessKeyID, "s3.* keys take precedence over rest.* aliases") | ||
| require.Equal(t, "SK", got.SecretAccessKey) | ||
|
|
||
| _, err = staticCredsFromProps(iceberg.Properties{keyRestAccessKeyID: "RAK"}) | ||
| require.ErrorIs(t, err, internalaws.ErrIncompleteStaticCredentials, "a lone rest.* access key must be an error") | ||
|
|
||
| _, err = staticCredsFromProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AK", | ||
| keyRestSecretAccessKey: "RSK", | ||
| }) | ||
| require.ErrorIs(t, err, internalaws.ErrIncompleteStaticCredentials, "a partial pair must not be completed with a field from the other namespace") | ||
|
|
||
| creds, err = staticCredsFromProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AK", | ||
| iceio.S3SecretAccessKey: "SK", | ||
| keyRestSessionToken: "RST", | ||
| }) | ||
| require.NoError(t, err) | ||
| got, err = creds.Retrieve(context.Background()) | ||
| require.NoError(t, err) | ||
| require.Equal(t, "AK", got.AccessKeyID) | ||
| require.Empty(t, got.SessionToken, "a complete s3.* pair must not inherit an unrelated rest.* session token") | ||
| } | ||
|
|
||
| // TestSigV4SignsWithPropsCredentials pins the wiring: the SigV4 Authorization | ||
| // header must be signed with the credentials carried in the catalog properties. | ||
| func TestSigV4SignsWithPropsCredentials(t *testing.T) { | ||
| var authHeader string | ||
| mux := http.NewServeMux() | ||
| mux.HandleFunc("/v1/config", func(w http.ResponseWriter, r *http.Request) { | ||
| json.NewEncoder(w).Encode(map[string]any{"defaults": map[string]any{}, "overrides": map[string]any{}}) | ||
| }) | ||
| mux.HandleFunc("/test", func(w http.ResponseWriter, r *http.Request) { | ||
| authHeader = r.Header.Get("Authorization") | ||
| w.WriteHeader(http.StatusOK) | ||
| }) | ||
| srv := httptest.NewServer(mux) | ||
| defer srv.Close() | ||
|
|
||
| cat, err := NewCatalog(context.Background(), "rest", srv.URL, | ||
| WithSigV4RegionSvc("us-east-1", "s3"), | ||
| WithAdditionalProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AKIDEXAMPLEPROPS", | ||
| iceio.S3SecretAccessKey: "secretexample", | ||
| })) | ||
| require.NoError(t, err) | ||
|
|
||
| req, err := http.NewRequestWithContext(context.Background(), http.MethodGet, srv.URL+"/test", nil) | ||
| require.NoError(t, err) | ||
| resp, err := cat.cl.Do(req) | ||
| require.NoError(t, err) | ||
| require.NoError(t, resp.Body.Close()) | ||
|
|
||
| require.Contains(t, authHeader, "Credential=AKIDEXAMPLEPROPS/", | ||
| "SigV4 must sign with the credentials from catalog properties, not the default chain") | ||
| } | ||
|
|
||
| // TestSigv4DoesNotSignCrossOriginRedirect pins that a redirect to a different | ||
| // origin is not re-signed, so the SigV4 Authorization header and session token | ||
| // never reach an unconfigured host. | ||
| func TestSigv4DoesNotSignCrossOriginRedirect(t *testing.T) { | ||
| var secondHit bool | ||
| var gotAuth, gotToken string | ||
| second := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { | ||
| secondHit = true | ||
| gotAuth = r.Header.Get("Authorization") | ||
| gotToken = r.Header.Get("X-Amz-Security-Token") | ||
| w.WriteHeader(http.StatusOK) | ||
| })) | ||
| defer second.Close() | ||
|
|
||
| var firstAuth string | ||
| mux := http.NewServeMux() | ||
| mux.HandleFunc("/v1/config", func(w http.ResponseWriter, r *http.Request) { | ||
| json.NewEncoder(w).Encode(map[string]any{"defaults": map[string]any{}, "overrides": map[string]any{}}) | ||
| }) | ||
| mux.HandleFunc("/redirect", func(w http.ResponseWriter, r *http.Request) { | ||
| firstAuth = r.Header.Get("Authorization") | ||
| http.Redirect(w, r, second.URL+"/landing", http.StatusTemporaryRedirect) | ||
| }) | ||
| first := httptest.NewServer(mux) | ||
| defer first.Close() | ||
|
|
||
| cat, err := NewCatalog(context.Background(), "rest", first.URL, | ||
| WithSigV4RegionSvc("us-east-1", "s3"), | ||
| WithAdditionalProps(iceberg.Properties{ | ||
| iceio.S3AccessKeyID: "AKIDEXAMPLEPROPS", | ||
| iceio.S3SecretAccessKey: "secretexample", | ||
| iceio.S3SessionToken: "SESSIONTOKENEXAMPLE", | ||
| })) | ||
| require.NoError(t, err) | ||
|
|
||
| req, err := http.NewRequestWithContext(context.Background(), http.MethodGet, first.URL+"/redirect", nil) | ||
| require.NoError(t, err) | ||
| resp, err := cat.cl.Do(req) | ||
| require.NoError(t, err) | ||
| require.NoError(t, resp.Body.Close()) | ||
|
|
||
| require.Contains(t, firstAuth, "Credential=AKIDEXAMPLEPROPS/", "the configured origin must still be signed") | ||
| require.True(t, secondHit, "the redirect target must be reached") | ||
| require.Empty(t, gotAuth, "the redirect target must not receive the SigV4 Authorization header") | ||
| require.Empty(t, gotToken, "the redirect target must not receive the session token") | ||
| } | ||
|
|
||
| func TestSplitIdentForPathRequiresNamespaceAndName(t *testing.T) { | ||
| cat := &Catalog{} | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -56,7 +56,7 @@ catalog: | |
| | `catalog.<name>.aws-profile` | AWS named profile for the Glue catalog. When unset, the AWS SDK default credential chain is used. | | ||
| | `catalog.<name>.sql-driver` | `database/sql` driver name for the SQL catalog. Maps to the `sql.driver` property. The default CLI binary only compiles in `sqliteshim`; other drivers require a custom build. | | ||
| | `catalog.<name>.sql-dialect` | SQL dialect for the SQL catalog (`postgres`, `mysql`, `sqlite`, `mssql`, `oracle`). Maps to the `sql.dialect` property. The default CLI binary only ships `sqlite` via `sqliteshim`; other dialects need a custom build with their drivers. | | ||
| | `catalog.<name>.rest.sigv4-enabled` | Enable AWS SigV4 signing for REST. | | ||
| | `catalog.<name>.rest.sigv4-enabled` | Enable AWS SigV4 signing for REST. When enabled, requests are signed with the `s3.*` credential properties if set (`s3.access-key-id` / `s3.secret-access-key` / `s3.session-token`), otherwise with the AWS default credential chain. The Java-client names (`rest.access-key-id` / `rest.secret-access-key` / `rest.session-token`) are accepted as aliases, resolved per field with the `s3.*` keys taking precedence. | | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Apart from the stale "resolved per field" wording (see options.go), this row is in the CLI YAML table, and the CLI config cannot set |
||
| | `catalog.<name>.rest.signing-name` | SigV4 service name. | | ||
| | `catalog.<name>.rest.signing-region` | SigV4 region. | | ||
|
|
||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.