zeroshade commented on code in PR #2001:
URL: https://github.com/apache/iceberg-go/pull/2001#discussion_r4158584452
##########
catalog/glue/glue.go:
##########
@@ -277,20 +283,61 @@ 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) {
+ // 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
+ }
+
+ var cfg catalog.CreateTableCfg
+ for _, opt := range opts {
+ opt(&cfg)
+ }
+
+ // S3 Tables manages storage, so an explicit location is rejected up
front
+ // (as pyiceberg does) rather than creating a table outside managed
storage.
+ if cfg.Location != "" {
+ federated, ferr := c.isS3TablesDatabase(ctx, database)
+ if ferr != nil {
+ return nil, ferr
+ }
+ if federated {
+ return nil, fmt.Errorf("cannot specify a location for
table %s.%s: S3 Tables manages storage automatically", database, tableName)
+ }
+ }
+
// 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...)
if err != nil {
- return nil, err
- }
+ // S3 Tables federated databases assign storage themselves, so
client-side
+ // location resolution fails with ErrNoDefaultLocation. Only
then probe for
+ // federation and retry via the S3 Tables path, keeping the
extra
+ // GetDatabase off every other create.
+ if !errors.Is(err, internal.ErrNoDefaultLocation) {
+ return nil, err
+ }
+ federated, ferr := c.isS3TablesDatabase(ctx, database)
+ if ferr != nil {
+ return nil, ferr
+ }
+ if !federated {
+ return nil, err
+ }
- database, tableName, err := identifierToGlueTable(identifier)
- if err != nil {
- return nil, err
+ return c.createS3TablesTable(ctx, database, tableName,
identifier, schema, opts...)
Review Comment:
This branch is only reached when `CreateStagedTable` fails with
`ErrNoDefaultLocation`. `getDefaultWarehouseLocation`
(catalog/internal/utils.go:221-231) falls back to the catalog `warehouse`
property, which is documented for Glue (website/src/configuration.md:53) and
reaches `c.props` through `catalog.Load` / `WithAwsProperties`. With
`warehouse` set, a create in an `aws:s3tables` database resolves
`<warehouse>/<db>.db/<tbl>`, never probes, writes metadata into the warehouse
bucket, and sends a full `EXTERNAL_TABLE` `TableInput` with that
`StorageDescriptor.Location` to the federated catalog. That is the same
outside-managed-storage create the explicit-location check at 304-312 rejects.
I traced this from the code and have not confirmed what S3 Tables does with
that request. pyiceberg avoids the problem by calling `_is_s3tables_database`
in `create_table` before any location resolution.
Fix: fetch the database once at the top of `CreateTable` (keeping the
`AccessDenied` tolerance for explicit-location creates), decide federation for
both default and explicit creates, and pass namespace props built from that
same result into `CreateStagedTable` through a closure (factor the conversion
out of `LoadNamespaceProperties`). That removes the second `GetDatabase` the S3
Tables path makes today (the tests mock it `.Times(2)`), keeps a single
round-trip for default creates, and lets you resolve the `CreateTableOpt`s
once, which closes my other thread. Please add a unit test with
`props{"warehouse": ...}` and a federated database that asserts the
minimal-entry create.
@tanmayrauth @laskoviymishka this replaces the lazy probe from 9791f9e. The
lazy probe misses any catalog with `warehouse` set, which is the implicit form
of the explicit-location bypass; reusing the one `GetDatabase` keeps the
round-trip saving the lazy probe was after.
##########
catalog/glue/glue.go:
##########
@@ -311,6 +358,160 @@ func (c *Catalog) CreateTable(ctx context.Context,
identifier table.Identifier,
return c.LoadTable(ctx, identifier)
}
+// isS3TablesDatabase reports whether the Glue database is federated to the
+// Amazon S3 Tables service, which owns table storage and location assignment.
+func (c *Catalog) isS3TablesDatabase(ctx context.Context, database string)
(bool, error) {
+ db, err := c.getDatabase(ctx, database)
+ if err != nil {
+ // Best-effort: a missing database or a caller lacking
glue:GetDatabase
+ // is treated as "not federated" so the generic create path can
proceed.
+ var apiErr smithy.APIError
+ if errors.Is(err, catalog.ErrNoSuchNamespace) ||
+ (errors.As(err, &apiErr) && apiErr.ErrorCode() ==
"AccessDeniedException") {
+ return false, nil
+ }
+
+ return false, err
+ }
+ if db == nil {
+ return false, nil
+ }
+
+ return isS3TablesFederatedDatabase(db), nil
+}
+
+func isS3TablesFederatedDatabase(db *types.Database) bool {
+ return db.FederatedDatabase != nil &&
+
strings.EqualFold(aws.ToString(db.FederatedDatabase.ConnectionType),
s3TablesConnectionType)
+}
+
+func isS3TablesFederatedTable(tbl *types.Table) bool {
+ return tbl.FederatedTable != nil &&
+
strings.EqualFold(aws.ToString(tbl.FederatedTable.ConnectionType),
s3TablesConnectionType)
+}
+
+// isS3TablesIcebergEntry reports whether a federated S3 Tables entry is
Iceberg,
+// accepting the minimal format=ICEBERG entry left before the metadata repoint
so
+// DropTable can clean it up after a failed create.
+func isS3TablesIcebergEntry(tbl *types.Table) bool {
+ if !isS3TablesFederatedTable(tbl) {
+ return false
+ }
+ tableType := tbl.Parameters[tableParamTableType]
+
+ return strings.EqualFold(tableType, glueTypeIceberg) ||
+ tableType == glueTypeIcebergRenaming ||
+ strings.EqualFold(tbl.Parameters[glueParamFormat],
glueTypeIceberg)
+}
+
+// createS3TablesTable creates a table in an S3 Tables federated database. The
+// service assigns storage, so a minimal entry is created first to allocate the
+// location, then updated with the written metadata pointer; on any later
+// failure the minimal entry is removed so no half-created table is left
behind.
+// On a commit failure the minimal Glue entry is rolled back and its metadata
+// object best-effort deleted; a stranded minimal entry after a double failure
is
+// still removable via DropTable, matching pyiceberg's direct delete_table.
+func (c *Catalog) createS3TablesTable(ctx context.Context, database, tableName
string, identifier table.Identifier, schema *iceberg.Schema, opts
...catalog.CreateTableOpt) (*table.Table, error) {
+ _, err := c.glueSvc.CreateTable(ctx, &glue.CreateTableInput{
+ CatalogId: c.catalogId,
+ DatabaseName: aws.String(database),
+ TableInput: &types.TableInput{
+ Name: aws.String(tableName),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ },
+ })
+ if err != nil {
+ return nil, fmt.Errorf("failed to allocate S3 Tables storage
for %s.%s: %w", database, tableName, err)
+ }
Review Comment:
An `AlreadyExistsException` from this allocate call comes back as `failed to
allocate S3 Tables storage ...` without `catalog.ErrTableAlreadyExists`. The
other create paths in this file map it (generic `CreateTable`, `RegisterTable`,
`CommitTable`), catalog/catalogtest/catalogtest.go:175 requires it, and
pyiceberg raises `TableAlreadyExistsError` at this same step. As written, a
create-if-not-exists caller that checks `errors.Is(err,
catalog.ErrTableAlreadyExists)` breaks on S3 Tables.
```suggestion
if err != nil {
if isAlreadyExistsException(err) {
return nil, fmt.Errorf("failed to create table %s.%s:
%w", database, tableName, catalog.ErrTableAlreadyExists)
}
return nil, fmt.Errorf("failed to allocate S3 Tables storage
for %s.%s: %w", database, tableName, err)
}
```
Please add a test where the allocate `CreateTable` returns
`AlreadyExistsException`, asserting `ErrorIs` and that `DeleteTable` is never
called.
##########
catalog/glue/glue.go:
##########
@@ -311,6 +358,160 @@ func (c *Catalog) CreateTable(ctx context.Context,
identifier table.Identifier,
return c.LoadTable(ctx, identifier)
}
+// isS3TablesDatabase reports whether the Glue database is federated to the
+// Amazon S3 Tables service, which owns table storage and location assignment.
+func (c *Catalog) isS3TablesDatabase(ctx context.Context, database string)
(bool, error) {
+ db, err := c.getDatabase(ctx, database)
+ if err != nil {
+ // Best-effort: a missing database or a caller lacking
glue:GetDatabase
+ // is treated as "not federated" so the generic create path can
proceed.
+ var apiErr smithy.APIError
+ if errors.Is(err, catalog.ErrNoSuchNamespace) ||
+ (errors.As(err, &apiErr) && apiErr.ErrorCode() ==
"AccessDeniedException") {
+ return false, nil
+ }
+
+ return false, err
+ }
+ if db == nil {
+ return false, nil
+ }
+
+ return isS3TablesFederatedDatabase(db), nil
+}
+
+func isS3TablesFederatedDatabase(db *types.Database) bool {
+ return db.FederatedDatabase != nil &&
+
strings.EqualFold(aws.ToString(db.FederatedDatabase.ConnectionType),
s3TablesConnectionType)
+}
+
+func isS3TablesFederatedTable(tbl *types.Table) bool {
+ return tbl.FederatedTable != nil &&
+
strings.EqualFold(aws.ToString(tbl.FederatedTable.ConnectionType),
s3TablesConnectionType)
+}
+
+// isS3TablesIcebergEntry reports whether a federated S3 Tables entry is
Iceberg,
+// accepting the minimal format=ICEBERG entry left before the metadata repoint
so
+// DropTable can clean it up after a failed create.
+func isS3TablesIcebergEntry(tbl *types.Table) bool {
+ if !isS3TablesFederatedTable(tbl) {
+ return false
+ }
+ tableType := tbl.Parameters[tableParamTableType]
+
+ return strings.EqualFold(tableType, glueTypeIceberg) ||
+ tableType == glueTypeIcebergRenaming ||
+ strings.EqualFold(tbl.Parameters[glueParamFormat],
glueTypeIceberg)
+}
+
+// createS3TablesTable creates a table in an S3 Tables federated database. The
+// service assigns storage, so a minimal entry is created first to allocate the
+// location, then updated with the written metadata pointer; on any later
+// failure the minimal entry is removed so no half-created table is left
behind.
+// On a commit failure the minimal Glue entry is rolled back and its metadata
+// object best-effort deleted; a stranded minimal entry after a double failure
is
+// still removable via DropTable, matching pyiceberg's direct delete_table.
+func (c *Catalog) createS3TablesTable(ctx context.Context, database, tableName
string, identifier table.Identifier, schema *iceberg.Schema, opts
...catalog.CreateTableOpt) (*table.Table, error) {
+ _, err := c.glueSvc.CreateTable(ctx, &glue.CreateTableInput{
+ CatalogId: c.catalogId,
+ DatabaseName: aws.String(database),
+ TableInput: &types.TableInput{
+ Name: aws.String(tableName),
+ Parameters: map[string]string{glueParamFormat:
glueTypeIceberg},
+ },
+ })
+ if err != nil {
+ return nil, fmt.Errorf("failed to allocate S3 Tables storage
for %s.%s: %w", database, tableName, err)
+ }
+
+ if err := c.commitS3TablesTable(ctx, database, tableName, identifier,
schema, opts...); err != nil {
+ // Roll back with a detached context so a cancelled create
still removes
+ // the minimal entry (S3 Tables enforces per-account table
limits).
+ cleanupCtx, cancel :=
context.WithTimeout(context.WithoutCancel(ctx), renameCleanupTimeout)
+ defer cancel()
+ if _, delErr := c.glueSvc.DeleteTable(cleanupCtx,
&glue.DeleteTableInput{
+ CatalogId: c.catalogId,
+ DatabaseName: aws.String(database),
+ Name: aws.String(tableName),
+ }); delErr != nil {
+ return nil, errors.Join(err, fmt.Errorf("failed to
clean up allocated table %s.%s: %w", database, tableName, delErr))
+ }
+
+ return nil, err
+ }
+
+ // The table is committed; load it outside the rollback scope so a
transient
+ // read failure does not delete an already-created table.
+ return c.LoadTable(ctx, identifier)
+}
+
+// commitS3TablesTable reads the service-assigned location, writes the Iceberg
+// metadata to it, and points the Glue entry at that metadata. It does not
reload
+// the table: the caller does that only after a successful commit, so a read
+// failure never triggers a rollback of an already-created table.
+func (c *Catalog) commitS3TablesTable(ctx context.Context, database, tableName
string, identifier table.Identifier, schema *iceberg.Schema, opts
...catalog.CreateTableOpt) error {
+ // Write managed-table metadata with the catalog's configured
credentials, not
+ // the ambient default chain, so the metadata PUT uses the same
principal as
+ // the Glue calls. Covers WriteMetadata and the staged.FS cleanup below.
+ ctx = utils.WithAwsConfig(ctx, c.awsCfg)
Review Comment:
Together with :340 and :584 this fixes the blocker from my 09-21 review. The
regression test I asked for is still missing: bc55241 only touches glue.go, and
every test here uses `file://`, so removing any one of these wraps would not
fail a test. `TestGlueRegisterTableAlreadyExists` already registers a scheme
backed by `iceio.NewMemFS()`; do the same with a factory that records
`utils.GetAwsConfig(ctx)`, and assert the recorded config is `c.awsCfg` for the
generic `CreateTable`, for `commitS3TablesTable` (including the `fs.Remove`
cleanup after a failed `UpdateTable`), and for `CommitTable`.
##########
catalog/glue/glue.go:
##########
@@ -866,7 +1073,12 @@ func (c *Catalog) getRawTable(ctx context.Context,
database, tableName string) (
return nil, fmt.Errorf("failed to get table %s.%s: missing Glue
table response", database, tableName)
}
- if aws.ToString(tblRes.Table.TableType) != glueTableType {
+ // A standard Iceberg table is an EXTERNAL_TABLE. An S3 Tables
federated entry
+ // instead carries the service's own TableType (e.g. "customer"), so
accept it
+ // when it is federated to S3 Tables and marked Iceberg — keeping the
+ // relaxation scoped to that case rather than every database.
+ isExternalTable := aws.ToString(tblRes.Table.TableType) == glueTableType
+ if !isExternalTable && !isS3TablesIcebergEntry(tblRes.Table) {
return nil, fmt.Errorf("table %s.%s is not an EXTERNAL_TABLE",
database, tableName)
}
Review Comment:
This gate also lets federated tables into `RenameTable` (glue.go:705-754,
outside the diff); before this PR they failed here with `is not an
EXTERNAL_TABLE`. `RenameTable` copies the source `StorageDescriptor.Location`
(the managed location) and `TableType` into a Glue `CreateTable` on the
destination, which is an explicit location of the kind `CreateTable` now
rejects, and then claims and deletes the source. If S3 Tables reclaims storage
on delete, as the rollback reasoning in 8ad630a assumes, the renamed entry can
end up pointing at reclaimed metadata. I have not confirmed the service
behaviour here. pyiceberg has the same unguarded path, but in this repo the
path is newly reachable and has no test. Either reject rename when
`isS3TablesFederatedTable(fromGlueTable)`, or add rename to the gated live test
and show the destination still loads after the source is gone.
##########
catalog/glue/glue_test.go:
##########
@@ -2607,3 +2616,586 @@ 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 TestGlueIsS3TablesDatabase(t *testing.T) {
+ tests := []struct {
+ name string
+ connectionType string
+ getErr error
+ want bool
+ wantErr bool
+ }{
+ {name: "federated to s3 tables", connectionType:
"aws:s3tables", want: true},
+ {name: "federated case insensitive", connectionType:
"AWS:S3Tables", want: true},
+ {name: "federated to another source", connectionType:
"aws:redshift", want: false},
+ {name: "not federated", connectionType: "", want: false},
+ {name: "missing database is not federated", getErr:
&types.EntityNotFoundException{}, want: false},
+ {name: "access denied is not federated", getErr:
&smithy.GenericAPIError{Code: "AccessDeniedException"}, want: false},
+ {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{}}
+ got, err :=
cat.isS3TablesDatabase(context.Background(), "test_database")
+ if tt.wantErr {
+ require.Error(t, err)
+ } else {
+ require.NoError(t, err)
+ require.Equal(t, tt.want, got)
+ }
+ mockGlueSvc.AssertExpectations(t)
+ })
+ }
+}
+
+// TestGlueCreateTableS3TablesFederated exercises the full two-phase create:
+// detect federation, allocate storage with a minimal entry, read the assigned
+// location, write metadata to it, and repoint the Glue entry.
+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).Times(2)
+
+ 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).Times(2)
+ 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).Times(2)
+ 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).Times(2)
+ 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 S3 Tables database, which manages
storage
+// itself; the table is never 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).Times(2)
+ 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)
+
+ 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)
Review Comment:
Nit: this test never calls `mockGlueSvc.AssertExpectations(t)`. More
broadly, no new test sets `catalogId`. The code passes `c.catalogId` on all
five new Glue calls (allocate `CreateTable`, `GetTable`, `UpdateTable`,
rollback `DeleteTable`, and the `GetDatabase` probe), but no test pins that,
and for a federated catalog it is the whole contract. One happy-path test with
`catalogId: aws.String("123456789012:s3tablescatalog/bucket")` and inputs
matched on it would cover it. The gated live test also stops at create and
reload; one commit (for example a property update) would cover the
`CommitTable` path these tables now reach.
--
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]