zeroshade commented on code in PR #2001:
URL: https://github.com/apache/iceberg-go/pull/2001#discussion_r4168186455
##########
catalog/glue/glue.go:
##########
@@ -277,20 +283,49 @@ var _ catalog.Closer = (*Catalog)(nil)
// This function will create the metadata file in S3 using the catalog and
table properties,
// to determine the bucket and key for the metadata location.
func (c *Catalog) CreateTable(ctx context.Context, identifier
table.Identifier, schema *iceberg.Schema, opts ...catalog.CreateTableOpt)
(*table.Table, error) {
- // The reporter is resolved once at construction (see NewCatalog), so a
bad
- // metrics-reporter-impl already failed there — no per-op guard is
needed
- // before mutating the catalog, and the trailing LoadTable reuses the
cached
- // reporter.
- staged, err := internal.CreateStagedTable(ctx, c.props,
c.LoadNamespaceProperties, identifier, schema, opts...)
+ // A missing namespace is reported before touching Glue, matching the
contract
+ // callers rely on (an identifier without a database is not a missing
table).
+ if len(identifier) < 2 {
+ return nil, fmt.Errorf("%w: missing namespace or invalid
identifier %v", catalog.ErrNoSuchNamespace, identifier)
+ }
+
+ database, tableName, err := identifierToGlueTable(identifier)
if err != nil {
return nil, err
}
- database, tableName, err := identifierToGlueTable(identifier)
+ var cfg catalog.CreateTableCfg
Review Comment:
Nit: `cfg` is only read in the federated branch, but this loop runs on every
create. Non-federated creates now evaluate the options twice, where `main`
evaluated them once. Moving the loop inside the `isS3TablesFederatedDatabase`
branch fixes that, and the glue tests still pass. Your pushback on the S3
Tables path holds: `WithLocation(managedLocation)` is appended last.
##########
catalog/internal/utils.go:
##########
@@ -40,6 +40,10 @@ import (
"github.com/klauspost/compress/zstd"
)
+// ErrNoDefaultLocation is returned when a table has no explicit location and
the
+// namespace and catalog define no default warehouse.
+var ErrNoDefaultLocation = errors.New("no default path set, please specify a
location when creating a table")
Review Comment:
Nit: nothing uses `ErrNoDefaultLocation` since the lazy probe was removed in
3d49411. Reverting this file would keep the PR Glue-only.
##########
catalog/glue/glue_test.go:
##########
@@ -2607,3 +2617,866 @@ func TestTableOperationsRejectEmptyIdentifiers(t
*testing.T) {
require.ErrorIs(t, err, catalog.ErrNoSuchTable)
}
}
+
+func s3TablesTestSchema() *iceberg.Schema {
+ return iceberg.NewSchemaWithIdentifiers(1, []int{1},
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.Int64Type{}, Required: true},
+ iceberg.NestedField{ID: 2, Name: "name", Type:
iceberg.StringType{}, Required: false},
+ )
+}
+
+func federatedDatabaseOutput(connectionType string) *glue.GetDatabaseOutput {
+ db := &types.Database{Name: aws.String("test_database")}
+ if connectionType != "" {
+ db.FederatedDatabase = &types.FederatedDatabase{ConnectionType:
aws.String(connectionType)}
+ }
+
+ return &glue.GetDatabaseOutput{Database: db}
+}
+
+func TestGlueLookupDatabase(t *testing.T) {
+ tests := []struct {
+ name string
+ connectionType string
+ getErr error
+ wantNil bool
+ wantFederated bool
+ wantErr bool
+ }{
+ {name: "federated to s3 tables", connectionType:
"aws:s3tables", wantFederated: true},
+ {name: "federated case insensitive", connectionType:
"AWS:S3Tables", wantFederated: true},
+ {name: "federated to another source", connectionType:
"aws:redshift"},
+ {name: "not federated", connectionType: ""},
+ {name: "missing database tolerated", getErr:
&types.EntityNotFoundException{}, wantNil: true},
+ {name: "access denied tolerated", getErr:
&smithy.GenericAPIError{Code: "AccessDeniedException"}, wantNil: true},
+ {name: "get database error", getErr: errors.New("boom"),
wantErr: true},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ mockGlueSvc := &mockGlueClient{}
+ if tt.getErr != nil {
+ mockGlueSvc.On("GetDatabase", mock.Anything,
&glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ },
mock.Anything).Return((*glue.GetDatabaseOutput)(nil), tt.getErr).Once()
+ } else {
+ mockGlueSvc.On("GetDatabase", mock.Anything,
&glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ },
mock.Anything).Return(federatedDatabaseOutput(tt.connectionType), nil).Once()
+ }
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg:
&aws.Config{}}
+ db, err := cat.lookupDatabase(context.Background(),
"test_database")
+ if tt.wantErr {
+ require.Error(t, err)
+ } else {
+ require.NoError(t, err)
+ if tt.wantNil {
+ require.Nil(t, db)
+ } else {
+ require.NotNil(t, db)
+ require.Equal(t, tt.wantFederated,
isS3TablesFederatedDatabase(db))
+ }
+ }
+ mockGlueSvc.AssertExpectations(t)
+ })
+ }
+}
+
+// TestGlueCreateTableS3TablesFederated exercises the full two-phase create.
+func TestGlueCreateTableS3TablesFederated(t *testing.T) {
+ ctx := context.Background()
+ managedLocation := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.MatchedBy(func(in
*glue.CreateTableInput) bool {
+ return aws.ToString(in.TableInput.Name) == "test_table" &&
+ in.TableInput.Parameters[glueParamFormat] ==
glueTypeIceberg &&
+ in.TableInput.StorageDescriptor == nil
+ }), mock.Anything).Return(&glue.CreateTableOutput{}, nil).Once()
+
+ allocated := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ TableType: aws.String("customer"),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: allocated},
nil).Once()
+
+ // LoadTable at the end reads this entry; UpdateTable's Run below fills
in the
+ // iceberg parameters (including the metadata pointer) before it is
read.
+ loaded := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ // S3 Tables reports its own service TableType plus a
FederatedTable marker,
+ // exercising the relaxed getRawTable gate on the create
success path.
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ var capturedMetadataLocation string
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.MatchedBy(func(in
*glue.UpdateTableInput) bool {
+ // The repoint sends EXTERNAL_TABLE even though the allocated
entry reports
+ // the service type "customer"; live testing confirmed S3
Tables accepts it.
+ return in.TableInput != nil && aws.ToString(in.VersionId) ==
"1" &&
+ in.TableInput.Parameters[tableParamTableType] ==
glueTypeIceberg &&
+ aws.ToString(in.TableInput.TableType) == glueTableType
+ }), mock.Anything).Run(func(args mock.Arguments) {
+ in := args.Get(1).(*glue.UpdateTableInput)
+ capturedMetadataLocation =
in.TableInput.Parameters[tableParamMetadataLocation]
+ loaded.Parameters[tableParamTableType] = glueTypeIceberg
+ loaded.Parameters[tableParamMetadataLocation] =
capturedMetadataLocation
+ }).Return(&glue.UpdateTableOutput{}, nil).Once()
+
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: loaded}, nil)
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ tbl, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.NoError(t, err)
+ require.Equal(t, TableIdentifier("test_database", "test_table"),
tbl.Identifier())
+ require.Equal(t, schema.Fields(), tbl.Schema().Fields())
+ require.Contains(t, tbl.MetadataLocation(), managedLocation)
+ require.Equal(t, capturedMetadataLocation, tbl.MetadataLocation())
+ require.FileExists(t, strings.TrimPrefix(capturedMetadataLocation,
"file://"))
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesCleanupOnFailure verifies the allocated entry is
+// deleted when the second phase fails, leaving no half-created table behind.
+func TestGlueCreateTableS3TablesCleanupOnFailure(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String("")},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "did not assign a storage location")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesCleanupErrorWrapped surfaces both the original
+// failure and the cleanup failure when deleting the allocated entry also
fails.
+func TestGlueCreateTableS3TablesCleanupErrorWrapped(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return((*glue.GetTableOutput)(nil), errors.New("get
boom")).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.DeleteTableOutput)(nil), errors.New("delete
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "get boom")
+ require.ErrorContains(t, err, "failed to clean up allocated table")
+ require.ErrorContains(t, err, "delete boom")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesAllocateError returns early without cleanup when
+// the initial allocation call itself fails.
+func TestGlueCreateTableS3TablesAllocateError(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.CreateTableOutput)(nil), errors.New("allocate
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "failed to allocate S3 Tables storage")
+ mockGlueSvc.AssertNotCalled(t, "GetTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesRejectsExplicitLocation verifies an explicit
location
+// is refused for a federated database and nothing is created.
+func TestGlueCreateTableS3TablesRejectsExplicitLocation(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema,
+ catalog.WithLocation("file:///tmp/whatever"))
+ require.ErrorContains(t, err, "S3 Tables manages storage automatically")
+ mockGlueSvc.AssertNotCalled(t, "CreateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableExplicitLocationNonFederated confirms an explicit
location
+// on a non-federated database still takes the generic path.
+func TestGlueCreateTableExplicitLocationNonFederated(t *testing.T) {
+ ctx := context.Background()
+ location := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput(""), nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ // The trailing reload is not the point here; fail it fast to avoid a
full
+ // metadata round-trip. What matters is that the generic create path
runs.
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return((*glue.GetTableOutput)(nil), errors.New("load
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema,
+ catalog.WithLocation(location))
+ require.ErrorContains(t, err, "load boom")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesMissingVersionId fails and rolls back when the
+// allocated entry has no Glue version id to commit against.
+func TestGlueCreateTableS3TablesMissingVersionId(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String("file:///tmp/whatever")},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "Glue table version id is missing")
+ mockGlueSvc.AssertNotCalled(t, "UpdateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableNonFederatedFallsThrough confirms a non-federated
database
+// with no location still hits the generic path (and its "no default path"
error).
+func TestGlueCreateTableNonFederatedFallsThrough(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput(""), nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "no default path set")
+ mockGlueSvc.AssertNotCalled(t, "CreateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesFederatedIntegration is gated by
TEST_S3TABLES_CATALOG_ID
+// (<account-id>:s3tablescatalog/<bucket>) and TEST_S3TABLES_DATABASE.
+func TestGlueCreateTableS3TablesFederatedIntegration(t *testing.T) {
+ catalogID := os.Getenv("TEST_S3TABLES_CATALOG_ID")
+ dbName := os.Getenv("TEST_S3TABLES_DATABASE")
+ if catalogID == "" || dbName == "" {
+ t.Skip()
+ }
+ assert := require.New(t)
+ ctx := context.Background()
+ awsCfg, err := config.LoadDefaultConfig(ctx)
+ assert.NoError(err)
+ ctlg, err := NewCatalog(WithAwsConfig(awsCfg),
WithAwsProperties(AwsProperties{CatalogIdKey: catalogID}))
+ assert.NoError(err)
+
+ tableName := fmt.Sprintf("it_%d", time.Now().UnixNano())
+ ident := TableIdentifier(dbName, tableName)
+ schema := s3TablesTestSchema()
+
+ tbl, err := ctlg.CreateTable(ctx, ident, schema)
+ assert.NoError(err)
+ defer func() { assert.NoError(ctlg.DropTable(ctx, ident)) }()
+
+ assert.Equal(ident, tbl.Identifier())
+ assert.Equal(schema.Fields(), tbl.Schema().Fields())
+ assert.Contains(tbl.MetadataLocation(), "--table-s3", "metadata must
land in the S3 Tables managed location")
+
+ reloaded, err := ctlg.LoadTable(ctx, ident)
+ assert.NoError(err)
+ assert.Equal(schema.Fields(), reloaded.Schema().Fields())
+ assert.Equal(tbl.MetadataLocation(), reloaded.MetadataLocation())
+}
+
+// TestGlueGetRawTableTableType covers the TableType gate: EXTERNAL_TABLE and
+// federated iceberg entries pass, anything else is rejected.
+func TestGlueGetRawTableTableType(t *testing.T) {
+ tests := []struct {
+ name string
+ tableType string
+ params map[string]string
+ federated bool
+ wantErr bool
+ }{
+ {
+ name: "standard external table",
+ tableType: glueTableType,
+ params: map[string]string{tableParamTableType:
glueTypeIceberg},
+ },
+ {
+ name: "s3 tables federated iceberg",
+ tableType: "customer",
+ params: map[string]string{tableParamTableType:
glueTypeIceberg},
+ federated: true,
+ },
+ {
+ name: "s3 tables federated renaming",
+ tableType: "customer",
+ params: map[string]string{tableParamTableType:
glueTypeIcebergRenaming},
+ federated: true,
+ },
+ {
+ name: "non-federated iceberg param unexpected type
rejected",
+ tableType: "customer",
+ params: map[string]string{tableParamTableType:
glueTypeIceberg},
+ wantErr: true,
+ },
+ {
+ name: "non-iceberg unexpected type rejected",
+ tableType: "VIRTUAL_VIEW",
+ params: map[string]string{},
+ wantErr: true,
+ },
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ mockGlueSvc := &mockGlueClient{}
+ glueTable := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ TableType: aws.String(tt.tableType),
+ Parameters: tt.params,
+ }
+ if tt.federated {
+ glueTable.FederatedTable =
&types.FederatedTable{ConnectionType: aws.String(s3TablesConnectionType)}
+ }
+ mockGlueSvc.On("GetTable", mock.Anything,
&glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table:
glueTable}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg:
&aws.Config{}}
+ tbl, err := cat.getRawTable(context.Background(),
"test_database", "test_table")
+ if tt.wantErr {
+ require.ErrorContains(t, err, "is not an
EXTERNAL_TABLE")
+ } else {
+ require.NoError(t, err)
+ require.Equal(t, tt.tableType,
aws.ToString(tbl.TableType))
+ }
+ mockGlueSvc.AssertExpectations(t)
+ })
+ }
+}
+
+// TestGlueCreateTableS3TablesRollbackOnMetadataWriteFailure verifies a failed
+// metadata write rolls back the allocated entry without issuing UpdateTable.
+func TestGlueCreateTableS3TablesRollbackOnMetadataWriteFailure(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ // An unregistered IO scheme fails deterministically at
metadata write,
+ // proving the rollback is triggered by the write and not by a
later step.
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String("unregisteredfs://bucket/table")},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "scheme not registered")
+ mockGlueSvc.AssertNotCalled(t, "UpdateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesRollbackOnUpdateFailure verifies a rejected
repoint
+// rolls back the allocated entry.
+func TestGlueCreateTableS3TablesRollbackOnUpdateFailure(t *testing.T) {
+ ctx := context.Background()
+ dir := t.TempDir()
+ managedLocation := "file://" + dir
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }}, nil).Once()
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.UpdateTableOutput)(nil), errors.New("update
rejected")).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "failed to commit S3 Tables table")
+ mockGlueSvc.AssertExpectations(t)
+
+ // The best-effort cleanup should have removed the metadata object it
wrote.
+ leftover, _ := filepath.Glob(filepath.Join(dir, "metadata",
"*.metadata.json"))
+ require.Empty(t, leftover)
+}
+
+// TestGlueCreateTableS3TablesNoRollbackOnLoadFailure verifies that a transient
+// failure of the final reload does not delete an already-committed table.
+func TestGlueCreateTableS3TablesNoRollbackOnLoadFailure(t *testing.T) {
+ ctx := context.Background()
+ managedLocation := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }}, nil).Once()
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.UpdateTableOutput{}, nil).Once()
+ // The trailing reload (LoadTable -> GetTable) fails transiently.
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return((*glue.GetTableOutput)(nil), errors.New("load
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "load boom")
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesRollbackDetachesContext verifies rollback still
runs
+// on a detached context when the create is cancelled.
+func TestGlueCreateTableS3TablesRollbackDetachesContext(t *testing.T) {
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, mock.Anything,
mock.Anything).
+ Return(federatedDatabaseOutput("aws:s3tables"), nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ // The commit load fails and the create is cancelled at the same moment.
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Run(func(mock.Arguments) { cancel() }).
+ Return((*glue.GetTableOutput)(nil), errors.New("commit
boom")).Once()
+ // The rollback must still see a live (detached) context, not the
cancelled one.
+ mockGlueSvc.On("DeleteTable",
+ mock.MatchedBy(func(c context.Context) bool { return c.Err() ==
nil }),
+ mock.Anything, mock.Anything).Return(&glue.DeleteTableOutput{},
nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "commit boom")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueDropTableRemovesStrandedS3TablesEntry verifies DropTable removes a
+// minimal entry (format=ICEBERG, no table_type) left by a failed create.
+func TestGlueDropTableRemovesStrandedS3TablesEntry(t *testing.T) {
+ ctx := context.Background()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ require.NoError(t, cat.DropTable(ctx, TableIdentifier("test_database",
"test_table")))
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// recordingMemFS registers an in-memory FileIO that records the context's AWS
+// config on every resolution.
+func recordingMemFS(t *testing.T, scheme string) (*iceio.MemFS,
*[]*aws.Config) {
+ t.Helper()
+ memFS := iceio.NewMemFS()
+ recorded := new([]*aws.Config)
+ iceio.Unregister(scheme)
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ *recorded = append(*recorded, utils.GetAwsConfig(ctx))
+
+ return memFS, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return memFS, recorded
+}
+
+func assertRecordedConfig(t *testing.T, want *aws.Config, recorded
*[]*aws.Config) {
+ t.Helper()
+ require.NotEmpty(t, *recorded, "expected the FileIO to be resolved at
least once")
+ for _, got := range *recorded {
+ require.Same(t, want, got, "metadata IO must use the catalog's
configured AWS config")
+ }
+}
+
+// TestGlueCreateTableS3TablesAllocateAlreadyExists verifies AlreadyExists
maps to
+// ErrTableAlreadyExists without rolling back a table this call did not create.
+func TestGlueCreateTableS3TablesAllocateAlreadyExists(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, &types.AlreadyExistsException{
+ Message: aws.String("Table already exists"),
+ }).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorIs(t, err, catalog.ErrTableAlreadyExists)
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesWarehouseSet verifies a `warehouse` property
does not
+// bypass federation: the create still goes through the minimal-entry allocate.
+func TestGlueCreateTableS3TablesWarehouseSet(t *testing.T) {
+ ctx := context.Background()
+ managedLocation := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ // The minimal allocate entry carries no StorageDescriptor; a
warehouse-located
+ // create would send one. Matching on that pins the two-phase path.
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.MatchedBy(func(in
*glue.CreateTableInput) bool {
+ return in.TableInput.Parameters[glueParamFormat] ==
glueTypeIceberg &&
+ in.TableInput.StorageDescriptor == nil
+ }), mock.Anything).Return(&glue.CreateTableOutput{}, nil).Once()
+
+ allocated := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ TableType: aws.String("customer"),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return(&glue.GetTableOutput{Table: allocated}, nil).Once()
+
+ loaded := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Run(func(args mock.Arguments) {
+ in := args.Get(1).(*glue.UpdateTableInput)
+ loaded.Parameters[tableParamTableType] = glueTypeIceberg
+ loaded.Parameters[tableParamMetadataLocation] =
in.TableInput.Parameters[tableParamMetadataLocation]
+ }).Return(&glue.UpdateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return(&glue.GetTableOutput{Table: loaded}, nil)
+
+ cat := &Catalog{
+ glueSvc: mockGlueSvc,
+ awsCfg: &aws.Config{},
+ props: iceberg.Properties{"warehouse":
"file:///tmp/warehouse"},
+ }
+ tbl, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.NoError(t, err)
+ require.Contains(t, tbl.MetadataLocation(), managedLocation)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesPassesCatalogID pins that every Glue call on the
+// two-phase create path carries the federated catalog id.
+func TestGlueCreateTableS3TablesPassesCatalogID(t *testing.T) {
+ ctx := context.Background()
+ managedLocation := "file://" + t.TempDir()
+ catalogID := "123456789012:s3tablescatalog/bucket"
+ schema := s3TablesTestSchema()
+
+ matchID := func(id *string) bool { return aws.ToString(id) == catalogID
}
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.MatchedBy(func(context.Context) bool
{ return true }),
+ mock.MatchedBy(func(in *glue.GetDatabaseInput) bool { return
matchID(in.CatalogId) }),
+ mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything,
+ mock.MatchedBy(func(in *glue.CreateTableInput) bool { return
matchID(in.CatalogId) }),
+ mock.Anything).Return(&glue.CreateTableOutput{}, nil).Once()
+
+ allocated := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("GetTable", mock.Anything,
+ mock.MatchedBy(func(in *glue.GetTableInput) bool { return
matchID(in.CatalogId) }),
+ mock.Anything).Return(&glue.GetTableOutput{Table: allocated},
nil).Once()
+
+ loaded := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("UpdateTable", mock.Anything,
+ mock.MatchedBy(func(in *glue.UpdateTableInput) bool { return
matchID(in.CatalogId) }),
+ mock.Anything).Run(func(args mock.Arguments) {
+ in := args.Get(1).(*glue.UpdateTableInput)
+ loaded.Parameters[tableParamTableType] = glueTypeIceberg
+ loaded.Parameters[tableParamMetadataLocation] =
in.TableInput.Parameters[tableParamMetadataLocation]
+ }).Return(&glue.UpdateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything,
+ mock.MatchedBy(func(in *glue.GetTableInput) bool { return
matchID(in.CatalogId) }),
+ mock.Anything).Return(&glue.GetTableOutput{Table: loaded}, nil)
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}, catalogId:
aws.String(catalogID)}
+ tbl, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.NoError(t, err)
+ require.Contains(t, tbl.MetadataLocation(), managedLocation)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueRenameTableS3TablesRejected verifies rename is refused for a
federated
+// S3 Tables table, which owns its managed storage, before any write.
+func TestGlueRenameTableS3TablesRejected(t *testing.T) {
+ ctx := context.Background()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("from_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("from_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{tableParamTableType:
glueTypeIceberg},
+ }}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.RenameTable(ctx,
+ TableIdentifier("test_database", "from_table"),
+ TableIdentifier("test_database", "to_table"))
+ require.ErrorContains(t, err, "renaming is not supported for S3 Tables
managed tables")
+ mockGlueSvc.AssertNotCalled(t, "CreateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableUsesCatalogAwsConfig pins that the generic create path
writes
+// and reads metadata with the catalog's configured AWS config.
+func TestGlueCreateTableUsesCatalogAwsConfig(t *testing.T) {
+ ctx := context.Background()
+ const scheme = "gluecreatecredcfg"
+ _, recorded := recordingMemFS(t, scheme)
+ awsCfg := &aws.Config{Region: "cred-regression"}
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, mock.Anything,
mock.Anything).
+ Return(federatedDatabaseOutput(""), nil).Once()
+ var metadataLoc string
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.MatchedBy(func(in
*glue.CreateTableInput) bool {
+ metadataLoc =
in.TableInput.Parameters[tableParamMetadataLocation]
+
+ return true
+ }), mock.Anything).Return(&glue.CreateTableOutput{}, nil).Once()
+
+ loaded := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ TableType: aws.String(glueTableType),
+ Parameters: map[string]string{},
+ }
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Run(func(mock.Arguments) {
+ loaded.Parameters[tableParamTableType] = glueTypeIceberg
+ loaded.Parameters[tableParamMetadataLocation] =
metadataLoc
+ }).Return(&glue.GetTableOutput{Table: loaded}, nil)
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: awsCfg, props:
iceberg.Properties{"warehouse": scheme + "://bucket"}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.NoError(t, err)
+ assertRecordedConfig(t, awsCfg, recorded)
+}
+
+// TestGlueCommitS3TablesTableUsesCatalogAwsConfig pins the catalog AWS config
on
+// the metadata write and the fs.Remove cleanup after a failed UpdateTable.
+func TestGlueCommitS3TablesTableUsesCatalogAwsConfig(t *testing.T) {
+ ctx := context.Background()
+ const scheme = "gluecommits3credcfg"
+ _, recorded := recordingMemFS(t, scheme)
+ managedLocation := scheme + "://bucket/test_table"
+ awsCfg := &aws.Config{Region: "cred-regression"}
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, mock.Anything,
mock.Anything).
+ Return(federatedDatabaseOutput("aws:s3tables"), nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }}, nil).Once()
+ // Fail the repoint so the fs.Remove cleanup runs; the entry is then
rolled back.
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.UpdateTableOutput)(nil), errors.New("update
boom")).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, mock.Anything,
mock.Anything).
Review Comment:
Nit: the rollback `DeleteTable` is the only new Glue call whose `CatalogId`
no test checks. Setting `catalogId` on this test's `Catalog` and matching it
here with `mock.MatchedBy` would cover it. With that change the test fails if
`CatalogId` is dropped from `glue.go:426`.
##########
catalog/glue/glue_test.go:
##########
@@ -2607,3 +2617,866 @@ func TestTableOperationsRejectEmptyIdentifiers(t
*testing.T) {
require.ErrorIs(t, err, catalog.ErrNoSuchTable)
}
}
+
+func s3TablesTestSchema() *iceberg.Schema {
+ return iceberg.NewSchemaWithIdentifiers(1, []int{1},
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.Int64Type{}, Required: true},
+ iceberg.NestedField{ID: 2, Name: "name", Type:
iceberg.StringType{}, Required: false},
+ )
+}
+
+func federatedDatabaseOutput(connectionType string) *glue.GetDatabaseOutput {
+ db := &types.Database{Name: aws.String("test_database")}
+ if connectionType != "" {
+ db.FederatedDatabase = &types.FederatedDatabase{ConnectionType:
aws.String(connectionType)}
+ }
+
+ return &glue.GetDatabaseOutput{Database: db}
+}
+
+func TestGlueLookupDatabase(t *testing.T) {
+ tests := []struct {
+ name string
+ connectionType string
+ getErr error
+ wantNil bool
+ wantFederated bool
+ wantErr bool
+ }{
+ {name: "federated to s3 tables", connectionType:
"aws:s3tables", wantFederated: true},
+ {name: "federated case insensitive", connectionType:
"AWS:S3Tables", wantFederated: true},
+ {name: "federated to another source", connectionType:
"aws:redshift"},
+ {name: "not federated", connectionType: ""},
+ {name: "missing database tolerated", getErr:
&types.EntityNotFoundException{}, wantNil: true},
+ {name: "access denied tolerated", getErr:
&smithy.GenericAPIError{Code: "AccessDeniedException"}, wantNil: true},
+ {name: "get database error", getErr: errors.New("boom"),
wantErr: true},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ mockGlueSvc := &mockGlueClient{}
+ if tt.getErr != nil {
+ mockGlueSvc.On("GetDatabase", mock.Anything,
&glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ },
mock.Anything).Return((*glue.GetDatabaseOutput)(nil), tt.getErr).Once()
+ } else {
+ mockGlueSvc.On("GetDatabase", mock.Anything,
&glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ },
mock.Anything).Return(federatedDatabaseOutput(tt.connectionType), nil).Once()
+ }
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg:
&aws.Config{}}
+ db, err := cat.lookupDatabase(context.Background(),
"test_database")
+ if tt.wantErr {
+ require.Error(t, err)
+ } else {
+ require.NoError(t, err)
+ if tt.wantNil {
+ require.Nil(t, db)
+ } else {
+ require.NotNil(t, db)
+ require.Equal(t, tt.wantFederated,
isS3TablesFederatedDatabase(db))
+ }
+ }
+ mockGlueSvc.AssertExpectations(t)
+ })
+ }
+}
+
+// TestGlueCreateTableS3TablesFederated exercises the full two-phase create.
+func TestGlueCreateTableS3TablesFederated(t *testing.T) {
+ ctx := context.Background()
+ managedLocation := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.MatchedBy(func(in
*glue.CreateTableInput) bool {
+ return aws.ToString(in.TableInput.Name) == "test_table" &&
+ in.TableInput.Parameters[glueParamFormat] ==
glueTypeIceberg &&
+ in.TableInput.StorageDescriptor == nil
+ }), mock.Anything).Return(&glue.CreateTableOutput{}, nil).Once()
+
+ allocated := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ VersionId: aws.String("1"),
+ TableType: aws.String("customer"),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: allocated},
nil).Once()
+
+ // LoadTable at the end reads this entry; UpdateTable's Run below fills
in the
+ // iceberg parameters (including the metadata pointer) before it is
read.
+ loaded := &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ // S3 Tables reports its own service TableType plus a
FederatedTable marker,
+ // exercising the relaxed getRawTable gate on the create
success path.
+ TableType: aws.String("customer"),
+ FederatedTable: &types.FederatedTable{ConnectionType:
aws.String(s3TablesConnectionType)},
+ Parameters: map[string]string{},
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String(managedLocation)},
+ }
+ var capturedMetadataLocation string
+ mockGlueSvc.On("UpdateTable", mock.Anything, mock.MatchedBy(func(in
*glue.UpdateTableInput) bool {
+ // The repoint sends EXTERNAL_TABLE even though the allocated
entry reports
+ // the service type "customer"; live testing confirmed S3
Tables accepts it.
+ return in.TableInput != nil && aws.ToString(in.VersionId) ==
"1" &&
+ in.TableInput.Parameters[tableParamTableType] ==
glueTypeIceberg &&
+ aws.ToString(in.TableInput.TableType) == glueTableType
+ }), mock.Anything).Run(func(args mock.Arguments) {
+ in := args.Get(1).(*glue.UpdateTableInput)
+ capturedMetadataLocation =
in.TableInput.Parameters[tableParamMetadataLocation]
+ loaded.Parameters[tableParamTableType] = glueTypeIceberg
+ loaded.Parameters[tableParamMetadataLocation] =
capturedMetadataLocation
+ }).Return(&glue.UpdateTableOutput{}, nil).Once()
+
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: loaded}, nil)
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ tbl, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.NoError(t, err)
+ require.Equal(t, TableIdentifier("test_database", "test_table"),
tbl.Identifier())
+ require.Equal(t, schema.Fields(), tbl.Schema().Fields())
+ require.Contains(t, tbl.MetadataLocation(), managedLocation)
+ require.Equal(t, capturedMetadataLocation, tbl.MetadataLocation())
+ require.FileExists(t, strings.TrimPrefix(capturedMetadataLocation,
"file://"))
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesCleanupOnFailure verifies the allocated entry is
+// deleted when the second phase fails, leaving no half-created table behind.
+func TestGlueCreateTableS3TablesCleanupOnFailure(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String("")},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "did not assign a storage location")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesCleanupErrorWrapped surfaces both the original
+// failure and the cleanup failure when deleting the allocated entry also
fails.
+func TestGlueCreateTableS3TablesCleanupErrorWrapped(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return((*glue.GetTableOutput)(nil), errors.New("get
boom")).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.DeleteTableOutput)(nil), errors.New("delete
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "get boom")
+ require.ErrorContains(t, err, "failed to clean up allocated table")
+ require.ErrorContains(t, err, "delete boom")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesAllocateError returns early without cleanup when
+// the initial allocation call itself fails.
+func TestGlueCreateTableS3TablesAllocateError(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return((*glue.CreateTableOutput)(nil), errors.New("allocate
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "failed to allocate S3 Tables storage")
+ mockGlueSvc.AssertNotCalled(t, "GetTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertNotCalled(t, "DeleteTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesRejectsExplicitLocation verifies an explicit
location
+// is refused for a federated database and nothing is created.
+func TestGlueCreateTableS3TablesRejectsExplicitLocation(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema,
+ catalog.WithLocation("file:///tmp/whatever"))
+ require.ErrorContains(t, err, "S3 Tables manages storage automatically")
+ mockGlueSvc.AssertNotCalled(t, "CreateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableExplicitLocationNonFederated confirms an explicit
location
+// on a non-federated database still takes the generic path.
+func TestGlueCreateTableExplicitLocationNonFederated(t *testing.T) {
+ ctx := context.Background()
+ location := "file://" + t.TempDir()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput(""), nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ // The trailing reload is not the point here; fail it fast to avoid a
full
+ // metadata round-trip. What matters is that the generic create path
runs.
+ mockGlueSvc.On("GetTable", mock.Anything, mock.Anything, mock.Anything).
+ Return((*glue.GetTableOutput)(nil), errors.New("load
boom")).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema,
+ catalog.WithLocation(location))
+ require.ErrorContains(t, err, "load boom")
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesMissingVersionId fails and rolls back when the
+// allocated entry has no Glue version id to commit against.
+func TestGlueCreateTableS3TablesMissingVersionId(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput("aws:s3tables"),
nil).Once()
+ mockGlueSvc.On("CreateTable", mock.Anything, mock.Anything,
mock.Anything).
+ Return(&glue.CreateTableOutput{}, nil).Once()
+ mockGlueSvc.On("GetTable", mock.Anything, &glue.GetTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.GetTableOutput{Table: &types.Table{
+ Name: aws.String("test_table"),
+ DatabaseName: aws.String("test_database"),
+ StorageDescriptor: &types.StorageDescriptor{Location:
aws.String("file:///tmp/whatever")},
+ }}, nil).Once()
+ mockGlueSvc.On("DeleteTable", mock.Anything, &glue.DeleteTableInput{
+ DatabaseName: aws.String("test_database"),
+ Name: aws.String("test_table"),
+ }, mock.Anything).Return(&glue.DeleteTableOutput{}, nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "Glue table version id is missing")
+ mockGlueSvc.AssertNotCalled(t, "UpdateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableNonFederatedFallsThrough confirms a non-federated
database
+// with no location still hits the generic path (and its "no default path"
error).
+func TestGlueCreateTableNonFederatedFallsThrough(t *testing.T) {
+ ctx := context.Background()
+ schema := s3TablesTestSchema()
+
+ mockGlueSvc := &mockGlueClient{}
+ mockGlueSvc.On("GetDatabase", mock.Anything, &glue.GetDatabaseInput{
+ Name: aws.String("test_database"),
+ }, mock.Anything).Return(federatedDatabaseOutput(""), nil).Once()
+
+ cat := &Catalog{glueSvc: mockGlueSvc, awsCfg: &aws.Config{}}
+ _, err := cat.CreateTable(ctx, TableIdentifier("test_database",
"test_table"), schema)
+ require.ErrorContains(t, err, "no default path set")
+ mockGlueSvc.AssertNotCalled(t, "CreateTable", mock.Anything,
mock.Anything, mock.Anything)
+ mockGlueSvc.AssertExpectations(t)
+}
+
+// TestGlueCreateTableS3TablesFederatedIntegration is gated by
TEST_S3TABLES_CATALOG_ID
+// (<account-id>:s3tablescatalog/<bucket>) and TEST_S3TABLES_DATABASE.
+func TestGlueCreateTableS3TablesFederatedIntegration(t *testing.T) {
+ catalogID := os.Getenv("TEST_S3TABLES_CATALOG_ID")
+ dbName := os.Getenv("TEST_S3TABLES_DATABASE")
+ if catalogID == "" || dbName == "" {
+ t.Skip()
+ }
+ assert := require.New(t)
+ ctx := context.Background()
+ awsCfg, err := config.LoadDefaultConfig(ctx)
+ assert.NoError(err)
+ ctlg, err := NewCatalog(WithAwsConfig(awsCfg),
WithAwsProperties(AwsProperties{CatalogIdKey: catalogID}))
+ assert.NoError(err)
+
+ tableName := fmt.Sprintf("it_%d", time.Now().UnixNano())
+ ident := TableIdentifier(dbName, tableName)
+ schema := s3TablesTestSchema()
+
+ tbl, err := ctlg.CreateTable(ctx, ident, schema)
+ assert.NoError(err)
+ defer func() { assert.NoError(ctlg.DropTable(ctx, ident)) }()
+
+ assert.Equal(ident, tbl.Identifier())
+ assert.Equal(schema.Fields(), tbl.Schema().Fields())
+ assert.Contains(tbl.MetadataLocation(), "--table-s3", "metadata must
land in the S3 Tables managed location")
+
+ reloaded, err := ctlg.LoadTable(ctx, ident)
+ assert.NoError(err)
+ assert.Equal(schema.Fields(), reloaded.Schema().Fields())
+ assert.Equal(tbl.MetadataLocation(), reloaded.MetadataLocation())
Review Comment:
Nit: `CommitTable` can now reach S3 Tables tables, but nothing exercises it
against the live service. One `SetProperties` commit after the reload would
cover it:
```suggestion
assert.Equal(tbl.MetadataLocation(), reloaded.MetadataLocation())
tx := reloaded.NewTransaction()
assert.NoError(tx.SetProperties(iceberg.Properties{"it.commit": "1"}))
committed, err := tx.Commit(ctx)
assert.NoError(err)
assert.NotEqual(reloaded.MetadataLocation(),
committed.MetadataLocation())
assert.Equal("1", committed.Properties()["it.commit"])
```
--
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]