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


##########
catalog/glue/glue.go:
##########
@@ -307,6 +351,155 @@ 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 {
+       allocated, err := c.glueSvc.GetTable(ctx, &glue.GetTableInput{
+               CatalogId:    c.catalogId,
+               DatabaseName: aws.String(database),
+               Name:         aws.String(tableName),
+       })
+       if err != nil {
+               return fmt.Errorf("failed to load allocated S3 Tables table 
%s.%s: %w", database, tableName, err)
+       }
+       if allocated == nil || allocated.Table == nil || 
allocated.Table.StorageDescriptor == nil {
+               return fmt.Errorf("S3 Tables did not return a storage 
descriptor for %s.%s", database, tableName)
+       }
+       managedLocation := 
aws.ToString(allocated.Table.StorageDescriptor.Location)
+       if managedLocation == "" {
+               return fmt.Errorf("S3 Tables did not assign a storage location 
for %s.%s", database, tableName)
+       }
+       if allocated.Table.VersionId == nil {
+               return fmt.Errorf("cannot commit table %s.%s: because Glue 
table version id is missing", database, tableName)
+       }
+
+       // Copy rather than append onto the caller's opts, whose backing array 
may
+       // have spare capacity we would otherwise clobber.
+       stagedOpts := make([]catalog.CreateTableOpt, len(opts), len(opts)+1)
+       copy(stagedOpts, opts)
+       stagedOpts = append(stagedOpts, catalog.WithLocation(managedLocation))
+
+       staged, err := internal.CreateStagedTable(ctx, c.props, 
c.LoadNamespaceProperties, identifier, schema, stagedOpts...)
+       if err != nil {
+               return err
+       }
+
+       if err := internal.WriteMetadata(ctx, staged.Table); err != nil {

Review Comment:
   This is the only metadata write in the file that doesn't carry the catalog's 
AWS config:
   
   ```go
   if err := internal.WriteMetadata(ctx, staged.Table); err != nil {
   ```
   
   Compare `:511` and `:1101`, which both wrap with `utils.WithAwsConfig(ctx, 
c.awsCfg)`, and `convertGlueToIceberg`, which injects `c.awsCfg` on the read 
path. With a bare `ctx` the S3 FileIO resolver falls back to 
`LoadDefaultConfig`, so a catalog constructed with `WithAwsConfig` (or `glue.*` 
static credentials) calls Glue as principal A and PUTs managed-table metadata 
as ambient principal B — or fails outright despite valid configured credentials.
   
   Neither existing test can catch this: the unit path is `file://`, and the 
gated live test supplies `config.LoadDefaultConfig` as the catalog config, so 
configured and ambient are indistinguishable.
   
   ```go
   storageCtx := utils.WithAwsConfig(ctx, c.awsCfg)
   ```
   
   Use that for `WriteMetadata` and for the cleanup `staged.FS` lookup. Please 
also audit the existing Glue CreateTable/CommitTable metadata writes, so a 
newly created S3 Tables table doesn't switch principals on its next commit. A 
regression where the catalog config differs from the ambient default would pin 
it.



##########
catalog/glue/glue.go:
##########
@@ -307,6 +351,155 @@ 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 {
+       allocated, err := c.glueSvc.GetTable(ctx, &glue.GetTableInput{
+               CatalogId:    c.catalogId,
+               DatabaseName: aws.String(database),
+               Name:         aws.String(tableName),
+       })
+       if err != nil {
+               return fmt.Errorf("failed to load allocated S3 Tables table 
%s.%s: %w", database, tableName, err)
+       }
+       if allocated == nil || allocated.Table == nil || 
allocated.Table.StorageDescriptor == nil {
+               return fmt.Errorf("S3 Tables did not return a storage 
descriptor for %s.%s", database, tableName)
+       }
+       managedLocation := 
aws.ToString(allocated.Table.StorageDescriptor.Location)
+       if managedLocation == "" {
+               return fmt.Errorf("S3 Tables did not assign a storage location 
for %s.%s", database, tableName)
+       }
+       if allocated.Table.VersionId == nil {
+               return fmt.Errorf("cannot commit table %s.%s: because Glue 
table version id is missing", database, tableName)
+       }
+
+       // Copy rather than append onto the caller's opts, whose backing array 
may
+       // have spare capacity we would otherwise clobber.
+       stagedOpts := make([]catalog.CreateTableOpt, len(opts), len(opts)+1)
+       copy(stagedOpts, opts)
+       stagedOpts = append(stagedOpts, catalog.WithLocation(managedLocation))
+
+       staged, err := internal.CreateStagedTable(ctx, c.props, 
c.LoadNamespaceProperties, identifier, schema, stagedOpts...)

Review Comment:
   Non-blocking: the `CreateTableOpt` functions are now evaluated two or three 
times — once for the location pre-check, once in initial staging, and again in 
S3 Tables commit staging. The built-ins are idempotent so this is benign today, 
but these are public options, and a legal stateful implementation could make 
the pre-check observe different configuration than staging, quietly bypassing 
the invariant the pre-check exists to enforce. Parsing once and passing the 
resolved config into staging would remove that surprise.



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

To unsubscribe, e-mail: [email protected]

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


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

Reply via email to