This is an automated email from the ASF dual-hosted git repository.
zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/terraform-provider-iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new a262c9d feat(provider): add AWS SigV4 request signing (#118)
a262c9d is described below
commit a262c9db7257b822fb0329cb06ef7690589f63e5
Author: Jiajia Li <[email protected]>
AuthorDate: Sat Oct 3 00:08:37 2026 +0800
feat(provider): add AWS SigV4 request signing (#118)
* feat(provider): add AWS SigV4 request signing
* Address review feedback on SigV4 signing
- Move the settings under a nested auth.sigv4 attribute. Setting it
enables signing, so sigv4_enabled is removed. Key pairing is checked
with stringvalidator.AlsoRequires, and keys may not be empty.
- Resolve the signing region for every credential source and fail when
none is set. Pass it to the AWS config load, so credential providers
that call STS (assume role, web identity) get it too.
- Reject partial static keys before anything is signed instead of
falling back to the ambient credential chain. Static keys without a
region only read the region from the AWS environment.
- While any part of auth is not known until apply, send no request
rather than signing with guessed settings.
- Under SigV4, set headers before signing so they are signed, send an
Authorization entry in headers as Original-Authorization, and share
http.DefaultTransport as the path without SigV4 does.
- Replace credential_process errors with a fixed message, because the
AWS SDK quotes output it cannot parse, secrets included.
- Tests: one stub-server helper, an isolated AWS environment (including
AWS_EC2_METADATA_DISABLED), provider-server tests for validation and
Configure, and a stub STS test for the region.
---
docs/index.md | 54 +++++
go.mod | 6 +-
internal/provider/provider.go | 218 +++++++++++++++++++-
internal/provider/provider_config_test.go | 209 +++++++++++++++++++
internal/provider/provider_sigv4_test.go | 327 ++++++++++++++++++++++++++++++
5 files changed, 809 insertions(+), 5 deletions(-)
diff --git a/docs/index.md b/docs/index.md
index 031047d..26d43df 100644
--- a/docs/index.md
+++ b/docs/index.md
@@ -45,7 +45,61 @@ Use Terraform to interact with Iceberg REST Catalog
instances.
### Optional
+- `auth` (Attributes) Authentication settings for the Iceberg REST catalog.
(see [below for nested schema](#nestedatt--auth))
- `headers` (Map of String, Sensitive) The headers to use for authentication.
- `token` (String, Sensitive) The token to use for authentication.
- `type` (String) The type of catalog. Use 'rest' for a plain REST catalog.
- `warehouse` (String) The warehouse to use for the Iceberg REST catalog. This
will be passed as `warehouse` property in the catalog properties.
+
+<a id="nestedatt--auth"></a>
+### Nested Schema for `auth`
+
+Optional:
+
+- `sigv4` (Attributes) Sign requests with AWS Signature Version 4, as catalogs
such as AWS Glue require. Setting this attribute enables signing. (see [below
for nested schema](#nestedatt--auth--sigv4))
+
+<a id="nestedatt--auth--sigv4"></a>
+### Nested Schema for `auth.sigv4`
+
+Optional:
+
+- `access_key_id` (String, Sensitive) Access key ID. When omitted, the
standard AWS credential chain (environment, shared config, instance role) is
used.
+- `region` (String) Signing region. When omitted, the region from the AWS
environment (`AWS_REGION`, shared config) is used.
+- `secret_access_key` (String, Sensitive) Secret access key. Must be set
together with `access_key_id`.
+- `session_token` (String, Sensitive) Session token for temporary (STS)
credentials. Requires `access_key_id` and `secret_access_key`.
+- `signing_name` (String) Signing service name (the credential-scope service).
Defaults to `execute-api`. Use `glue` for AWS Glue.
+
+## AWS SigV4 authentication
+
+Some REST catalogs authenticate requests with AWS Signature Version 4 instead
of
+a bearer token. Setting `auth.sigv4` signs every request. The signing region
+comes from `region` or, when that is omitted, from the AWS environment; the
+provider reports an error when neither sets one. The signing name defaults to
+`execute-api`; override it for catalogs that scope signatures to a different
+service.
+
+Credentials come from `access_key_id` and `secret_access_key`, plus
+`session_token` for temporary credentials, when they are set, and from the
+standard AWS credential chain otherwise. Values from `headers` are set before
+signing, so they are part of the signature.
+
+When `token`, or an `Authorization` entry in `headers`, is also set, SigV4
keeps
+the `Authorization` header and that value is sent as `Original-Authorization`,
+matching the Java client, for catalogs that sit behind a SigV4 gateway and
+authenticate with OAuth themselves.
+
+```terraform
+# AWS Glue REST catalog
+provider "iceberg" {
+ catalog_uri = "https://glue.us-east-1.amazonaws.com/iceberg"
+ warehouse = "123456789012"
+
+ auth = {
+ sigv4 = {
+ region = "us-east-1"
+ signing_name = "glue"
+ # Omit the keys to use the standard AWS credential chain.
+ }
+ }
+}
+```
diff --git a/go.mod b/go.mod
index ef53fda..51f4a7b 100644
--- a/go.mod
+++ b/go.mod
@@ -19,6 +19,9 @@ go 1.25.8
require (
github.com/apache/iceberg-go v0.6.0
+ github.com/aws/aws-sdk-go-v2 v1.41.7
+ github.com/aws/aws-sdk-go-v2/config v1.32.17
+ github.com/aws/aws-sdk-go-v2/credentials v1.19.16
github.com/hashicorp/terraform-plugin-framework v1.19.0
github.com/hashicorp/terraform-plugin-framework-validators v0.19.0
github.com/hashicorp/terraform-plugin-go v0.31.0
@@ -39,9 +42,6 @@ require (
github.com/apache/arrow-go/v18 v18.6.0 // indirect
github.com/apache/thrift v0.24.0 // indirect
github.com/apparentlymart/go-textseg/v15 v15.0.0 // indirect
- github.com/aws/aws-sdk-go-v2 v1.41.7 // indirect
- github.com/aws/aws-sdk-go-v2/config v1.32.17 // indirect
- github.com/aws/aws-sdk-go-v2/credentials v1.19.16 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.23 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.23 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.23 // indirect
diff --git a/internal/provider/provider.go b/internal/provider/provider.go
index 7dddad7..9bd8dee 100644
--- a/internal/provider/provider.go
+++ b/internal/provider/provider.go
@@ -17,15 +17,25 @@ package provider
import (
"context"
+ "errors"
+ "fmt"
"net/http"
"github.com/apache/iceberg-go/catalog"
"github.com/apache/iceberg-go/catalog/rest"
+ "github.com/aws/aws-sdk-go-v2/aws"
+ "github.com/aws/aws-sdk-go-v2/config"
+ "github.com/aws/aws-sdk-go-v2/credentials"
+ "github.com/aws/aws-sdk-go-v2/credentials/processcreds"
+
"github.com/hashicorp/terraform-plugin-framework-validators/stringvalidator"
"github.com/hashicorp/terraform-plugin-framework/datasource"
+ "github.com/hashicorp/terraform-plugin-framework/path"
"github.com/hashicorp/terraform-plugin-framework/provider"
"github.com/hashicorp/terraform-plugin-framework/provider/schema"
"github.com/hashicorp/terraform-plugin-framework/resource"
+ "github.com/hashicorp/terraform-plugin-framework/schema/validator"
"github.com/hashicorp/terraform-plugin-framework/types"
+ "github.com/hashicorp/terraform-plugin-framework/types/basetypes"
)
var _ provider.Provider = &icebergProvider{}
@@ -44,6 +54,18 @@ type icebergProvider struct {
token string
warehouse string
headers map[string]string
+ sigv4 *sigv4Config
+}
+
+// sigv4Config holds the settings for signing requests with AWS SigV4.
+type sigv4Config struct {
+ region string
+ signingName string
+ accessKeyID string
+ secretAccessKey string
+ sessionToken string
+ // unknown is set while planning when auth depends on values known only
after apply.
+ unknown bool
}
// icebergProviderModel maps provider schema data to a Go type.
@@ -53,6 +75,19 @@ type icebergProviderModel struct {
Token types.String `tfsdk:"token"`
Warehouse types.String `tfsdk:"warehouse"`
Headers types.Map `tfsdk:"headers"`
+ Auth types.Object `tfsdk:"auth"`
+}
+
+type icebergAuthModel struct {
+ SigV4 types.Object `tfsdk:"sigv4"`
+}
+
+type icebergSigV4Model struct {
+ Region types.String `tfsdk:"region"`
+ SigningName types.String `tfsdk:"signing_name"`
+ AccessKeyID types.String `tfsdk:"access_key_id"`
+ SecretAccessKey types.String `tfsdk:"secret_access_key"`
+ SessionToken types.String `tfsdk:"session_token"`
}
// Metadata returns the provider type name.
@@ -88,6 +123,56 @@ func (p *icebergProvider) Schema(_ context.Context, _
provider.SchemaRequest, re
Sensitive: true,
ElementType: types.StringType,
},
+ "auth": schema.SingleNestedAttribute{
+ Description: "Authentication settings for the
Iceberg REST catalog.",
+ Optional: true,
+ Attributes: map[string]schema.Attribute{
+ "sigv4": schema.SingleNestedAttribute{
+ Description: "Sign requests
with AWS Signature Version 4, as catalogs such as AWS Glue require. Setting
this attribute enables signing.",
+ Optional: true,
+ Attributes:
map[string]schema.Attribute{
+ "region":
schema.StringAttribute{
+ Description:
"Signing region. When omitted, the region from the AWS environment
(`AWS_REGION`, shared config) is used.",
+ Optional:
true,
+ },
+ "signing_name":
schema.StringAttribute{
+ Description:
"Signing service name (the credential-scope service). Defaults to
`execute-api`. Use `glue` for AWS Glue.",
+ Optional:
true,
+ },
+ "access_key_id":
schema.StringAttribute{
+ Description:
"Access key ID. When omitted, the standard AWS credential chain (environment,
shared config, instance role) is used.",
+ Optional:
true,
+ Sensitive:
true,
+ Validators:
[]validator.String{
+
stringvalidator.LengthAtLeast(1),
+
stringvalidator.AlsoRequires(path.MatchRelative().AtParent().AtName("secret_access_key")),
+ },
+ },
+ "secret_access_key":
schema.StringAttribute{
+ Description:
"Secret access key. Must be set together with `access_key_id`.",
+ Optional:
true,
+ Sensitive:
true,
+ Validators:
[]validator.String{
+
stringvalidator.LengthAtLeast(1),
+
stringvalidator.AlsoRequires(path.MatchRelative().AtParent().AtName("access_key_id")),
+ },
+ },
+ "session_token":
schema.StringAttribute{
+ Description:
"Session token for temporary (STS) credentials. Requires `access_key_id` and
`secret_access_key`.",
+ Optional:
true,
+ Sensitive:
true,
+ Validators:
[]validator.String{
+
stringvalidator.LengthAtLeast(1),
+
stringvalidator.AlsoRequires(
+
path.MatchRelative().AtParent().AtName("access_key_id"),
+
path.MatchRelative().AtParent().AtName("secret_access_key"),
+ ),
+ },
+ },
+ },
+ },
+ },
+ },
},
}
}
@@ -144,13 +229,50 @@ func (p *icebergProvider) Configure(ctx context.Context,
req provider.ConfigureR
p.headers = headers
}
+ p.sigv4 = nil
+ if !data.Auth.IsNull() {
+ authValue, err := data.Auth.ToTerraformValue(ctx)
+ if err != nil {
+ resp.Diagnostics.AddError("Invalid auth configuration",
err.Error())
+
+ return
+ }
+
+ if !authValue.IsFullyKnown() {
+ // Guessing any part would sign with the wrong scope or
identity.
+ p.sigv4 = &sigv4Config{unknown: true}
+ } else {
+ var auth icebergAuthModel
+ resp.Diagnostics.Append(data.Auth.As(ctx, &auth,
basetypes.ObjectAsOptions{})...)
+ if resp.Diagnostics.HasError() {
+ return
+ }
+
+ if !auth.SigV4.IsNull() {
+ var m icebergSigV4Model
+ resp.Diagnostics.Append(auth.SigV4.As(ctx, &m,
basetypes.ObjectAsOptions{})...)
+ if resp.Diagnostics.HasError() {
+ return
+ }
+
+ p.sigv4 = &sigv4Config{
+ region: m.Region.ValueString(),
+ signingName:
m.SigningName.ValueString(),
+ accessKeyID:
m.AccessKeyID.ValueString(),
+ secretAccessKey:
m.SecretAccessKey.ValueString(),
+ sessionToken:
m.SessionToken.ValueString(),
+ }
+ }
+ }
+ }
+
resp.DataSourceData = p
resp.ResourceData = p
}
func (p *icebergProvider) NewCatalog(ctx context.Context) (catalog.Catalog,
error) {
opts := make([]rest.Option, 0)
- if p.token != "" {
+ if p.token != "" && p.sigv4 == nil {
opts = append(opts, rest.WithOAuthToken(p.token))
}
@@ -158,11 +280,103 @@ func (p *icebergProvider) NewCatalog(ctx
context.Context) (catalog.Catalog, erro
opts = append(opts, rest.WithWarehouseLocation(p.warehouse))
}
- opts = append(opts,
rest.WithCustomTransport(&headerRoundTripper{headers: p.headers}))
+ if p.sigv4 == nil {
+ opts = append(opts,
rest.WithCustomTransport(&headerRoundTripper{headers: p.headers}))
+ } else {
+ sigv4Opts, err := p.sigv4.options(ctx, p.token, p.headers)
+ if err != nil {
+ return nil, err
+ }
+ opts = append(opts, sigv4Opts...)
+ }
return rest.NewCatalog(ctx, p.catalogType, p.catalogURI, opts...)
}
+// options returns the iceberg-go options that sign requests with SigV4.
+// iceberg-go signs before a custom transport runs, so the configured headers
+// go through it to be signed with the request. SigV4 owns Authorization, so an
+// Authorization header or bearer token travels as Original-Authorization, as
+// the Java client does.
+func (s *sigv4Config) options(ctx context.Context, token string, headers
map[string]string) ([]rest.Option, error) {
+ cfg, err := s.awsConfig(ctx)
+ if err != nil {
+ return nil, err
+ }
+
+ signed := make(map[string]string, len(headers)+1)
+ for name, value := range headers {
+ if http.CanonicalHeaderKey(name) == "Authorization" {
+ name = "Original-Authorization"
+ }
+ signed[name] = value
+ }
+ if token != "" {
+ signed["Original-Authorization"] = "Bearer " + token
+ }
+
+ return []rest.Option{
+ rest.WithHeaders(signed),
+ rest.WithAwsConfig(cfg),
+ rest.WithSigV4RegionSvc(cfg.Region, s.signingName),
+ // Without a transport, iceberg-go builds one per catalog and
never
+ // closes its idle connections.
+ rest.WithCustomTransport(http.DefaultTransport),
+ }, nil
+}
+
+// awsConfig resolves the credentials and region to sign with. Static keys
+// replace the AWS credential chain. The region falls back to the AWS
+// environment, and credential providers that call STS get it as well.
+func (s *sigv4Config) awsConfig(ctx context.Context) (aws.Config, error) {
+ if s.unknown {
+ return aws.Config{}, errors.New("auth.sigv4 is not known until
apply, so requests cannot be signed yet")
+ }
+ static := s.accessKeyID != "" && s.secretAccessKey != ""
+ if !static && (s.accessKeyID != "" || s.secretAccessKey != "" ||
s.sessionToken != "") {
+ // Never fall back to the AWS credential chain when keys were
configured.
+ return aws.Config{}, errors.New("auth.sigv4: access_key_id and
secret_access_key must be set together, and session_token requires both")
+ }
+
+ var load []func(*config.LoadOptions) error
+ if static {
+ creds :=
credentials.NewStaticCredentialsProvider(s.accessKeyID, s.secretAccessKey,
s.sessionToken)
+ if s.region != "" {
+ return aws.Config{Region: s.region, Credentials:
creds}, nil
+ }
+ // Only the region comes from the AWS environment, not its
credential chain.
+ load = append(load, config.WithCredentialsProvider(creds))
+ }
+ if s.region != "" {
+ load = append(load, config.WithRegion(s.region))
+ }
+
+ cfg, err := config.LoadDefaultConfig(ctx, load...)
+ if err != nil {
+ return aws.Config{}, fmt.Errorf("load the AWS configuration for
auth.sigv4: %w", err)
+ }
+ if cfg.Region == "" {
+ return aws.Config{}, errors.New("auth.sigv4.region is required:
no region is configured in the AWS environment")
+ }
+ cfg.Credentials = processOutputHidden{cfg.Credentials}
+
+ return cfg, nil
+}
+
+// processOutputHidden keeps credential_process output out of errors: the AWS
+// SDK quotes output it cannot parse, and that output carries the secrets.
+type processOutputHidden struct{ aws.CredentialsProvider }
+
+func (p processOutputHidden) Retrieve(ctx context.Context) (aws.Credentials,
error) {
+ creds, err := p.CredentialsProvider.Retrieve(ctx)
+ var processErr *processcreds.ProviderError
+ if errors.As(err, &processErr) {
+ return creds, errors.New("credential_process failed; its output
is not shown because it can contain secrets")
+ }
+
+ return creds, err
+}
+
type headerRoundTripper struct {
headers map[string]string
}
diff --git a/internal/provider/provider_config_test.go
b/internal/provider/provider_config_test.go
new file mode 100644
index 0000000..c91b236
--- /dev/null
+++ b/internal/provider/provider_config_test.go
@@ -0,0 +1,209 @@
+// 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 provider
+
+import (
+ "context"
+ "maps"
+ "strings"
+ "testing"
+
+ "github.com/hashicorp/terraform-plugin-framework/provider"
+ "github.com/hashicorp/terraform-plugin-framework/providerserver"
+ "github.com/hashicorp/terraform-plugin-go/tfprotov6"
+ "github.com/hashicorp/terraform-plugin-go/tftypes"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// schemaTypes returns the Terraform types of the provider config, auth and
auth.sigv4.
+func schemaTypes(t *testing.T) (config, auth, sigv4 tftypes.Object) {
+ t.Helper()
+ ctx := context.Background()
+
+ var schemaResp provider.SchemaResponse
+ New()().Schema(ctx, provider.SchemaRequest{}, &schemaResp)
+ config, ok :=
schemaResp.Schema.Type().TerraformType(ctx).(tftypes.Object)
+ require.True(t, ok)
+ auth, ok = config.AttributeTypes["auth"].(tftypes.Object)
+ require.True(t, ok)
+ sigv4, ok = auth.AttributeTypes["sigv4"].(tftypes.Object)
+ require.True(t, ok)
+
+ return config, auth, sigv4
+}
+
+// object builds a value of typ with the given attributes set and the rest
null.
+func object(typ tftypes.Object, set map[string]tftypes.Value) tftypes.Value {
+ vals := make(map[string]tftypes.Value, len(typ.AttributeTypes))
+ for name, attrType := range typ.AttributeTypes {
+ vals[name] = tftypes.NewValue(attrType, nil)
+ }
+ maps.Copy(vals, set)
+
+ return tftypes.NewValue(typ, vals)
+}
+
+// authProviderConfig builds a provider configuration with catalog_uri and the
given auth value.
+func authProviderConfig(t *testing.T, auth tftypes.Value)
*tfprotov6.DynamicValue {
+ t.Helper()
+ configType, _, _ := schemaTypes(t)
+ config, err := tfprotov6.NewDynamicValue(configType, object(configType,
map[string]tftypes.Value{
+ "catalog_uri": str("http://localhost:8181"),
+ "auth": auth,
+ }))
+ require.NoError(t, err)
+
+ return &config
+}
+
+// sigv4ProviderConfig builds a provider configuration that sets the given
auth.sigv4 attributes.
+func sigv4ProviderConfig(t *testing.T, sigv4 map[string]tftypes.Value)
*tfprotov6.DynamicValue {
+ t.Helper()
+ _, authType, sigv4Type := schemaTypes(t)
+
+ return authProviderConfig(t, object(authType,
map[string]tftypes.Value{"sigv4": object(sigv4Type, sigv4)}))
+}
+
+func errorDiagnostics(diags []*tfprotov6.Diagnostic) []string {
+ var errs []string
+ for _, d := range diags {
+ if d.Severity == tfprotov6.DiagnosticSeverityError {
+ errs = append(errs, d.Summary+": "+d.Detail)
+ }
+ }
+
+ return errs
+}
+
+func str(v string) tftypes.Value { return tftypes.NewValue(tftypes.String, v) }
+
+func TestConfigureReadsSigV4Settings(t *testing.T) {
+ p := &icebergProvider{}
+ server, err := providerserver.NewProtocol6WithError(p)()
+ require.NoError(t, err)
+
+ resp, err := server.ConfigureProvider(context.Background(),
&tfprotov6.ConfigureProviderRequest{
+ Config: sigv4ProviderConfig(t, map[string]tftypes.Value{
+ "region": str("us-east-1"),
+ "signing_name": str("glue"),
+ "access_key_id": str("AKID"),
+ "secret_access_key": str("SECRET"),
+ "session_token": str("TOKEN"),
+ }),
+ })
+ require.NoError(t, err)
+ require.Empty(t, errorDiagnostics(resp.Diagnostics))
+
+ assert.Equal(t, &sigv4Config{
+ region: "us-east-1",
+ signingName: "glue",
+ accessKeyID: "AKID",
+ secretAccessKey: "SECRET",
+ sessionToken: "TOKEN",
+ }, p.sigv4)
+}
+
+func TestConfigureClearsSigV4WhenRemoved(t *testing.T) {
+ _, authType, _ := schemaTypes(t)
+ p := &icebergProvider{}
+ server, err := providerserver.NewProtocol6WithError(p)()
+ require.NoError(t, err)
+
+ for _, config := range []*tfprotov6.DynamicValue{
+ sigv4ProviderConfig(t, map[string]tftypes.Value{"region":
str("us-east-1")}),
+ authProviderConfig(t, tftypes.NewValue(authType, nil)),
+ } {
+ resp, err := server.ConfigureProvider(context.Background(),
&tfprotov6.ConfigureProviderRequest{Config: config})
+ require.NoError(t, err)
+ require.Empty(t, errorDiagnostics(resp.Diagnostics))
+ }
+
+ assert.Nil(t, p.sigv4, "a provider configured again without auth must
stop signing")
+}
+
+func TestValidateSigV4RequiresCompleteKeyPair(t *testing.T) {
+ for name, tc := range map[string]struct {
+ sigv4 map[string]tftypes.Value
+ want []string
+ }{
+ "access key only":
{map[string]tftypes.Value{"access_key_id": str("AKID")},
[]string{`"auth.sigv4.secret_access_key" must be specified`}},
+ "secret key only":
{map[string]tftypes.Value{"secret_access_key": str("SECRET")},
[]string{`"auth.sigv4.access_key_id" must be specified`}},
+ "session token only":
{map[string]tftypes.Value{"session_token": str("TOKEN")},
[]string{`"auth.sigv4.access_key_id" must be specified`,
`"auth.sigv4.secret_access_key" must be specified`}},
+ "complete keys":
{map[string]tftypes.Value{"access_key_id": str("AKID"), "secret_access_key":
str("SECRET"), "session_token": str("TOKEN")}, nil},
+ "empty keys": {map[string]tftypes.Value{"access_key_id":
str(""), "secret_access_key": str("")}, []string{
+ `auth.sigv4.access_key_id string length must be at
least 1`, `auth.sigv4.secret_access_key string length must be at least 1`,
+ }},
+ "empty session token":
{map[string]tftypes.Value{"access_key_id": str("AKID"), "secret_access_key":
str("SECRET"), "session_token": str("")}, []string{
+ `auth.sigv4.session_token string length must be at
least 1`,
+ }},
+ "ambient chain": {map[string]tftypes.Value{"region":
str("us-east-1")}, nil},
+ "secret unknown at plan": {map[string]tftypes.Value{
+ "access_key_id": str("AKID"),
+ "secret_access_key": tftypes.NewValue(tftypes.String,
tftypes.UnknownValue),
+ }, nil},
+ } {
+ t.Run(name, func(t *testing.T) {
+ server, err :=
providerserver.NewProtocol6WithError(New()())()
+ require.NoError(t, err)
+ resp, err :=
server.ValidateProviderConfig(context.Background(),
+
&tfprotov6.ValidateProviderConfigRequest{Config: sigv4ProviderConfig(t,
tc.sigv4)})
+ require.NoError(t, err)
+
+ errs := errorDiagnostics(resp.Diagnostics)
+ if len(tc.want) == 0 {
+ assert.Empty(t, errs)
+
+ return
+ }
+ for _, want := range tc.want {
+ assert.Contains(t, strings.Join(errs, "\n"),
want)
+ }
+ })
+ }
+}
+
+func TestNewCatalogRefusesSigV4SettingsUnknownAtPlan(t *testing.T) {
+ isolateAWSEnv(t)
+ _, authType, sigv4Type := schemaTypes(t)
+ unknown := tftypes.NewValue(tftypes.String, tftypes.UnknownValue)
+ withSigV4 := func(set map[string]tftypes.Value) tftypes.Value {
+ return object(authType, map[string]tftypes.Value{"sigv4":
object(sigv4Type, set)})
+ }
+
+ for name, auth := range map[string]tftypes.Value{
+ "keys": withSigV4(map[string]tftypes.Value{"region":
str("us-east-1"), "access_key_id": unknown, "secret_access_key": unknown}),
+ "region": withSigV4(map[string]tftypes.Value{"region":
unknown}),
+ "signing name": withSigV4(map[string]tftypes.Value{"region":
str("us-east-1"), "signing_name": unknown}),
+ "sigv4": object(authType,
map[string]tftypes.Value{"sigv4": tftypes.NewValue(sigv4Type,
tftypes.UnknownValue)}),
+ "auth": tftypes.NewValue(authType,
tftypes.UnknownValue),
+ } {
+ t.Run(name, func(t *testing.T) {
+ p := &icebergProvider{}
+ server, err := providerserver.NewProtocol6WithError(p)()
+ require.NoError(t, err)
+ resp, err :=
server.ConfigureProvider(context.Background(),
+ &tfprotov6.ConfigureProviderRequest{Config:
authProviderConfig(t, auth)})
+ require.NoError(t, err)
+ require.Empty(t, errorDiagnostics(resp.Diagnostics))
+
+ headers, err := configHeaders(t, p)
+ require.ErrorContains(t, err, "not known until apply")
+
+ assert.Nil(t, headers, "nothing may be sent until
auth.sigv4 is known")
+ })
+ }
+}
diff --git a/internal/provider/provider_sigv4_test.go
b/internal/provider/provider_sigv4_test.go
new file mode 100644
index 0000000..1ffb3c2
--- /dev/null
+++ b/internal/provider/provider_sigv4_test.go
@@ -0,0 +1,327 @@
+// 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 provider
+
+import (
+ "context"
+ "net"
+ "net/http"
+ "net/http/httptest"
+ "os"
+ "path/filepath"
+ "strings"
+ "sync/atomic"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// configHeaders creates a catalog against a stub REST server and returns the
+// headers of its /v1/config request, or nil when no request was sent.
+func configHeaders(t *testing.T, p *icebergProvider) (http.Header, error) {
+ t.Helper()
+ var headers http.Header
+ srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter,
r *http.Request) {
+ if r.URL.Path == "/v1/config" {
+ headers = r.Header.Clone()
+ }
+ w.Header().Set("Content-Type", "application/json")
+ _, _ = w.Write([]byte(`{"defaults":{},"overrides":{}}`))
+ }))
+ defer srv.Close()
+
+ p.catalogURI, p.catalogType = srv.URL, "rest"
+ _, err := p.NewCatalog(context.Background())
+
+ return headers, err
+}
+
+// isolateAWSEnv hides the host's AWS configuration and gives the AWS
+// credential chain a known ambient key pair.
+func isolateAWSEnv(t *testing.T) {
+ t.Helper()
+ for _, kv := range os.Environ() {
+ if name, _, _ := strings.Cut(kv, "="); strings.HasPrefix(name,
"AWS_") {
+ t.Setenv(name, "")
+ }
+ }
+ t.Setenv("AWS_CONFIG_FILE", "/nonexistent/config")
+ t.Setenv("AWS_SHARED_CREDENTIALS_FILE", "/nonexistent/credentials")
+ t.Setenv("AWS_EC2_METADATA_DISABLED", "true")
+ t.Setenv("AWS_ACCESS_KEY_ID", "AMBIENTKEY")
+ t.Setenv("AWS_SECRET_ACCESS_KEY", "AMBIENTSECRET")
+}
+
+func TestNewCatalogSignsRequestsWithSigV4(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{sigv4: &sigv4Config{
+ region: "us-east-1",
+ signingName: "glue",
+ accessKeyID: "AKIDEXAMPLE",
+ secretAccessKey: "SECRETEXAMPLE",
+ }})
+ require.NoError(t, err)
+
+ assert.Regexp(t, `^AWS4-HMAC-SHA256
Credential=AKIDEXAMPLE/\d{8}/us-east-1/glue/aws4_request,`,
+ headers.Get("Authorization"), "the credential scope carries the
configured keys, region and signing name")
+ assert.NotEmpty(t, headers.Get("x-amz-content-sha256"),
+ "SigV4 requires the x-amz-content-sha256 header")
+}
+
+func TestNewCatalogDefaultsSigningNameToExecuteAPI(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{sigv4: &sigv4Config{
+ region: "us-east-1",
+ accessKeyID: "AKIDEXAMPLE",
+ secretAccessKey: "SECRETEXAMPLE",
+ }})
+ require.NoError(t, err)
+
+ assert.Contains(t, headers.Get("Authorization"),
"/us-east-1/execute-api/aws4_request",
+ "an omitted signing name must default to execute-api")
+}
+
+func TestNewCatalogSignsWithSessionToken(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{sigv4: &sigv4Config{
+ region: "us-east-1",
+ accessKeyID: "AKIDEXAMPLE",
+ secretAccessKey: "SECRETEXAMPLE",
+ sessionToken: "SESSIONTOKEN",
+ }})
+ require.NoError(t, err)
+
+ assert.Equal(t, "SESSIONTOKEN", headers.Get("X-Amz-Security-Token"),
+ "temporary credentials must send X-Amz-Security-Token")
+}
+
+func TestNewCatalogWithoutSigV4SendsNoSignature(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{})
+ require.NoError(t, err)
+
+ assert.Empty(t, headers.Get("Authorization"), "plain config must not
send an Authorization header")
+ assert.Empty(t, headers.Get("x-amz-content-sha256"), "plain config must
not send x-amz-content-sha256")
+}
+
+func TestNewCatalogRelocatesBearerTokenUnderSigV4(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{
+ token: "bearer-example",
+ sigv4: &sigv4Config{region: "us-east-1", accessKeyID:
"AKIDEXAMPLE", secretAccessKey: "SECRETEXAMPLE"},
+ })
+ require.NoError(t, err)
+
+ assert.Contains(t, headers.Get("Authorization"), "AWS4-HMAC-SHA256",
+ "SigV4 owns the Authorization header")
+ assert.Equal(t, "Bearer bearer-example",
headers.Get("Original-Authorization"),
+ "the bearer token moves to Original-Authorization, as the Java
client does")
+ assert.Contains(t, headers.Get("Authorization"),
"original-authorization",
+ "the relocated header is part of the signed headers")
+}
+
+func TestNewCatalogRelocatesAuthorizationHeaderUnderSigV4(t *testing.T) {
+ for _, name := range []string{"Authorization", "authorization"} {
+ t.Run(name, func(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{
+ headers: map[string]string{name: "Bearer
from-headers"},
+ sigv4: &sigv4Config{region: "us-east-1",
accessKeyID: "AKIDEXAMPLE", secretAccessKey: "SECRETEXAMPLE"},
+ })
+ require.NoError(t, err)
+
+ assert.Equal(t, "Bearer from-headers",
headers.Get("Original-Authorization"),
+ "an Authorization header from headers is
relocated, as the Java client does")
+ assert.Contains(t, headers.Get("Authorization"),
"AWS4-HMAC-SHA256")
+ })
+ }
+}
+
+func TestNewCatalogSignsAmbientCredentialsWithConfiguredRegion(t *testing.T) {
+ isolateAWSEnv(t)
+ t.Setenv("AWS_REGION", "eu-west-1")
+
+ headers, err := configHeaders(t, &icebergProvider{sigv4:
&sigv4Config{region: "ap-southeast-2"}})
+ require.NoError(t, err)
+
+ assert.Regexp(t,
`Credential=AMBIENTKEY/\d{8}/ap-southeast-2/execute-api/aws4_request`,
+ headers.Get("Authorization"), "the configured region wins over
the AWS environment")
+}
+
+func TestNewCatalogResolvesRegionFromEnvironment(t *testing.T) {
+ isolateAWSEnv(t)
+ t.Setenv("AWS_REGION", "eu-west-1")
+
+ for name, tc := range map[string]struct {
+ sigv4 *sigv4Config
+ key string
+ }{
+ "static credentials": {&sigv4Config{accessKeyID:
"AKIDEXAMPLE", secretAccessKey: "SECRETEXAMPLE"}, "AKIDEXAMPLE"},
+ "ambient credentials": {&sigv4Config{}, "AMBIENTKEY"},
+ } {
+ t.Run(name, func(t *testing.T) {
+ headers, err := configHeaders(t,
&icebergProvider{sigv4: tc.sigv4})
+ require.NoError(t, err)
+
+ assert.Regexp(t,
`Credential=`+tc.key+`/\d{8}/eu-west-1/execute-api/aws4_request`,
+ headers.Get("Authorization"), "an omitted
region falls back to the AWS environment")
+ })
+ }
+}
+
+func TestNewCatalogRequiresARegion(t *testing.T) {
+ isolateAWSEnv(t)
+
+ for name, sigv4 := range map[string]*sigv4Config{
+ "static credentials": {accessKeyID: "AKIDEXAMPLE",
secretAccessKey: "SECRETEXAMPLE"},
+ "ambient credentials": {},
+ } {
+ t.Run(name, func(t *testing.T) {
+ headers, err := configHeaders(t,
&icebergProvider{sigv4: sigv4})
+ require.ErrorContains(t, err, "auth.sigv4.region is
required")
+
+ assert.Nil(t, headers, "no request may be signed with
an empty region")
+ })
+ }
+}
+
+func TestNewCatalogRejectsPartialStaticCredentials(t *testing.T) {
+ isolateAWSEnv(t)
+
+ for name, sigv4 := range map[string]*sigv4Config{
+ "access key only": {region: "us-east-1", accessKeyID:
"EXPLICITKEY"},
+ "secret key only": {region: "us-east-1", secretAccessKey:
"EXPLICITSECRET"},
+ "session token only": {region: "us-east-1", sessionToken:
"SESSIONTOKEN"},
+ } {
+ t.Run(name, func(t *testing.T) {
+ headers, err := configHeaders(t,
&icebergProvider{sigv4: sigv4})
+ require.ErrorContains(t, err, "must be set together")
+
+ assert.Nil(t, headers, "partial keys must not fall back
to the AWS credential chain")
+ })
+ }
+}
+
+func TestNewCatalogSignsCustomHeaders(t *testing.T) {
+ headers, err := configHeaders(t, &icebergProvider{
+ headers: map[string]string{"X-Iceberg-Access-Delegation":
"remote-signing", "X-Custom": "custom"},
+ sigv4: &sigv4Config{region: "us-east-1", accessKeyID:
"AKIDEXAMPLE", secretAccessKey: "SECRETEXAMPLE"},
+ })
+ require.NoError(t, err)
+
+ assert.Equal(t, []string{"remote-signing"},
headers.Values("X-Iceberg-Access-Delegation"),
+ "a configured header replaces the client default instead of
adding a second, unsigned value")
+ assert.Contains(t, headers.Get("Authorization"), "x-custom",
+ "configured headers are part of the signed headers")
+}
+
+func TestNewCatalogPassesConfiguredRegionToAmbientCredentials(t *testing.T) {
+ isolateAWSEnv(t)
+ t.Setenv("AWS_ACCESS_KEY_ID", "")
+ t.Setenv("AWS_SECRET_ACCESS_KEY", "")
+
+ // An assume-role profile: its credential provider signs an STS request
with
+ // the region it was built with.
+ var stsAuth string
+ sts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter,
r *http.Request) {
+ stsAuth = r.Header.Get("Authorization")
+ w.Header().Set("Content-Type", "text/xml")
+ _, _ = w.Write([]byte(`<AssumeRoleResponse
xmlns="https://sts.amazonaws.com/doc/2011-06-15/"><AssumeRoleResult><Credentials>`
+
+
`<AccessKeyId>STSKEY</AccessKeyId><SecretAccessKey>STSSECRET</SecretAccessKey><SessionToken>STSTOKEN</SessionToken>`
+
+
`<Expiration>2099-01-01T00:00:00Z</Expiration></Credentials></AssumeRoleResult></AssumeRoleResponse>`))
+ }))
+ defer sts.Close()
+ configFile := filepath.Join(t.TempDir(), "config")
+ require.NoError(t, os.WriteFile(configFile, []byte(`[profile assume]
+role_arn = arn:aws:iam::123456789012:role/example
+source_profile = base
+
+[profile base]
+aws_access_key_id = BASEKEY
+aws_secret_access_key = BASESECRET
+`), 0o600))
+ t.Setenv("AWS_CONFIG_FILE", configFile)
+ t.Setenv("AWS_PROFILE", "assume")
+ t.Setenv("AWS_ENDPOINT_URL_STS", sts.URL)
+
+ headers, err := configHeaders(t, &icebergProvider{sigv4:
&sigv4Config{region: "ap-southeast-2"}})
+ require.NoError(t, err)
+
+ assert.Contains(t, stsAuth, "/ap-southeast-2/sts/aws4_request",
+ "credential providers that call STS need the configured region
too")
+ assert.Regexp(t,
`Credential=STSKEY/\d{8}/ap-southeast-2/execute-api/aws4_request`,
headers.Get("Authorization"))
+}
+
+func TestNewCatalogResolvesRegionWithoutAmbientCredentialChain(t *testing.T) {
+ isolateAWSEnv(t)
+ t.Setenv("AWS_ACCESS_KEY_ID", "")
+ t.Setenv("AWS_SECRET_ACCESS_KEY", "")
+ configFile := filepath.Join(t.TempDir(), "config")
+ require.NoError(t, os.WriteFile(configFile, []byte(`[profile mfa]
+region = eu-central-1
+role_arn = arn:aws:iam::123456789012:role/example
+source_profile = base
+mfa_serial = arn:aws:iam::123456789012:mfa/example
+
+[profile base]
+aws_access_key_id = BASEKEY
+aws_secret_access_key = BASESECRET
+`), 0o600))
+ t.Setenv("AWS_CONFIG_FILE", configFile)
+ t.Setenv("AWS_PROFILE", "mfa")
+
+ headers, err := configHeaders(t, &icebergProvider{sigv4:
&sigv4Config{accessKeyID: "AKIDEXAMPLE", secretAccessKey: "SECRETEXAMPLE"}})
+ require.NoError(t, err, "static keys must not need the profile's MFA
credential chain")
+
+ assert.Regexp(t,
`Credential=AKIDEXAMPLE/\d{8}/eu-central-1/execute-api/aws4_request`,
headers.Get("Authorization"),
+ "the region still comes from the profile")
+}
+
+func TestNewCatalogHidesCredentialProcessOutput(t *testing.T) {
+ isolateAWSEnv(t)
+ t.Setenv("AWS_ACCESS_KEY_ID", "")
+ t.Setenv("AWS_SECRET_ACCESS_KEY", "")
+ configFile := filepath.Join(t.TempDir(), "config")
+ // The AWS SDK quotes output it cannot parse, and the output carries
secrets.
+ require.NoError(t, os.WriteFile(configFile, []byte("[profile
process]\ncredential_process = echo SECRETVALUE\n"), 0o600))
+ t.Setenv("AWS_CONFIG_FILE", configFile)
+ t.Setenv("AWS_PROFILE", "process")
+
+ _, err := configHeaders(t, &icebergProvider{sigv4: &sigv4Config{region:
"us-east-1"}})
+ require.ErrorContains(t, err, "credential_process")
+
+ assert.NotContains(t, err.Error(), "SECRETVALUE")
+}
+
+func TestNewCatalogSharesConnectionsUnderSigV4(t *testing.T) {
+ var conns atomic.Int32
+ srv := httptest.NewUnstartedServer(http.HandlerFunc(func(w
http.ResponseWriter, r *http.Request) {
+ w.Header().Set("Content-Type", "application/json")
+ _, _ = w.Write([]byte(`{"defaults":{},"overrides":{}}`))
+ }))
+ srv.Config.ConnState = func(_ net.Conn, state http.ConnState) {
+ if state == http.StateNew {
+ conns.Add(1)
+ }
+ }
+ srv.Start()
+ defer srv.Close()
+
+ for range 5 {
+ p := &icebergProvider{catalogURI: srv.URL, catalogType: "rest",
sigv4: &sigv4Config{
+ region: "us-east-1", accessKeyID: "AKIDEXAMPLE",
secretAccessKey: "SECRETEXAMPLE",
+ }}
+ _, err := p.NewCatalog(context.Background())
+ require.NoError(t, err)
+ }
+
+ assert.Equal(t, int32(1), conns.Load(), "catalogs reuse one connection
pool instead of leaving one open each")
+}