This is an automated email from the ASF dual-hosted git repository.
roryqi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 8e9ca0009d [#13565] fix: Validate table and column field lengths
(#13551)
8e9ca0009d is described below
commit 8e9ca0009d68dbf9edb5e313d3d51688b9c4cff1
Author: roryqi <[email protected]>
AuthorDate: Mon Sep 28 23:22:57 2026 +0800
[#13565] fix: Validate table and column field lengths (#13551)
### What changes were proposed in this pull request?
- Limit table and top-level column names to 128 characters.
- Limit top-level column comments to 4096 characters.
- Validate Gravitino table create and alter requests before catalog
changes.
- Validate Iceberg REST create and update requests before catalog
changes.
- Validate staged-create table names before committing the external
table.
- Validate registered Iceberg table metadata before register/overwrite
mutates the external catalog.
### Why are the changes needed?
The relational entity store limits table and column names to 128
characters and column comments to 4096 characters, but these constraints
were not enforced consistently in the logic layer.
Oversized values could reach the external catalog before failing while
writing Gravitino metadata, resulting in database-specific errors and
inconsistent metadata.
Fix: #13565
### Does this PR introduce _any_ user-facing change?
Yes. Oversized table names, top-level column names, and top-level column
comments are rejected with clear HTTP 400 errors.
Nested Iceberg fields are not restricted by the top-level column
metadata limits.
### How was this patch tested?
- `./gradlew :core:test --tests
org.apache.gravitino.meta.TestEntityFieldLimits --tests
org.apache.gravitino.catalog.TestTableOperationDispatcher`
- `./gradlew :iceberg:iceberg-common:test --tests
org.apache.gravitino.iceberg.common.ops.TestIcebergCatalogWrapper`
- `./gradlew :iceberg:iceberg-rest-server:test --tests
org.apache.gravitino.iceberg.service.dispatcher.TestIcebergTableOperationExecutor
--tests
org.apache.gravitino.iceberg.service.dispatcher.TestIcebergNamespaceOperationExecutor`
- `./gradlew spotlessApply`
---
.../org/apache/gravitino/EntityFieldLimits.java | 3 +
.../catalog/TableOperationDispatcher.java | 33 ++++
.../org/apache/gravitino/meta/ColumnEntity.java | 7 +-
.../org/apache/gravitino/meta/TableEntity.java | 4 +-
.../catalog/TestTableOperationDispatcher.java | 174 +++++++++++++++++++++
.../gravitino/meta/TestEntityFieldLimits.java | 32 ++++
.../iceberg/common/ops/IcebergCatalogWrapper.java | 15 ++
.../common/ops/TestIcebergCatalogWrapper.java | 26 +++
.../dispatcher/IcebergColumnFieldValidator.java | 49 ++++++
.../IcebergNamespaceOperationExecutor.java | 15 +-
.../dispatcher/IcebergTableOperationExecutor.java | 10 ++
.../service/dispatcher/TestIcebergAsyncPurge.java | 4 +
.../TestIcebergNamespaceOperationExecutor.java | 141 +++++++++++++++++
.../TestIcebergTableOperationExecutor.java | 139 ++++++++++++++++
.../service/rest/CatalogWrapperForTest.java | 25 +--
15 files changed, 662 insertions(+), 15 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/EntityFieldLimits.java
b/core/src/main/java/org/apache/gravitino/EntityFieldLimits.java
index 8c7edab14c..90077d5138 100644
--- a/core/src/main/java/org/apache/gravitino/EntityFieldLimits.java
+++ b/core/src/main/java/org/apache/gravitino/EntityFieldLimits.java
@@ -36,6 +36,9 @@ public final class EntityFieldLimits {
/** The maximum number of characters of an entity comment stored in a
256-character column. */
public static final int MAX_COMMENT_LENGTH = 256;
+ /** The maximum number of characters of a column comment. */
+ public static final int MAX_COLUMN_COMMENT_LENGTH = 4096;
+
private EntityFieldLimits() {}
/**
diff --git
a/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
b/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
index f4bfb96378..afcef4b385 100644
---
a/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
+++
b/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
@@ -18,6 +18,7 @@
*/
package org.apache.gravitino.catalog;
+import static org.apache.gravitino.Entity.EntityType.COLUMN;
import static org.apache.gravitino.Entity.EntityType.TABLE;
import static org.apache.gravitino.catalog.CapabilityHelpers.applyCapabilities;
import static
org.apache.gravitino.catalog.PropertiesMetadataHelpers.validatePropertyForCreate;
@@ -221,6 +222,9 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
Index[] indexes)
throws NoSuchSchemaException, TableAlreadyExistsException {
+ TableEntity.NAME.validate(ident.name(), TABLE);
+ validateColumns(columns);
+
// Load the schema to make sure the schema exists.
SchemaDispatcher schemaDispatcher = getSchemaDispatcher();
NameIdentifier schemaIdent = NameIdentifier.of(ident.namespace().levels());
@@ -275,9 +279,11 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
// to on the schema.
NameIdentifier nameIdentifierForLock = ident;
String schemaName = ident.namespace().level(2);
+
Arrays.stream(changes).forEach(TableOperationDispatcher::validateColumnChange);
for (TableChange change : changes) {
if (change instanceof TableChange.RenameTable) {
TableChange.RenameTable rename = (TableChange.RenameTable) change;
+ TableEntity.NAME.validate(rename.getNewName(), TABLE);
if (rename.getNewSchemaName().isPresent()
&& !rename.getNewSchemaName().get().equals(schemaName)) {
nameIdentifierForLock = getCatalogIdentifier(ident);
@@ -1055,6 +1061,33 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
combinedTable.tableFromGravitino().id()));
}
+ private static void validateColumns(Column[] columns) {
+ for (Column column : columns) {
+ ColumnEntity.NAME.validate(column.name(), COLUMN);
+ ColumnEntity.COMMENT.validate(column.comment(), COLUMN);
+ }
+ }
+
+ private static void validateColumnChange(TableChange change) {
+ if (change instanceof TableChange.AddColumn) {
+ TableChange.AddColumn addColumn = (TableChange.AddColumn) change;
+ if (addColumn.getFieldName().length == 1) {
+ ColumnEntity.NAME.validate(addColumn.getFieldName()[0], COLUMN);
+ ColumnEntity.COMMENT.validate(addColumn.getComment(), COLUMN);
+ }
+ } else if (change instanceof TableChange.RenameColumn) {
+ TableChange.RenameColumn renameColumn = (TableChange.RenameColumn)
change;
+ if (renameColumn.getFieldName().length == 1) {
+ ColumnEntity.NAME.validate(renameColumn.getNewName(), COLUMN);
+ }
+ } else if (change instanceof TableChange.UpdateColumnComment) {
+ TableChange.UpdateColumnComment updateComment =
(TableChange.UpdateColumnComment) change;
+ if (updateComment.getFieldName().length == 1) {
+ ColumnEntity.COMMENT.validate(updateComment.getNewComment(), COLUMN);
+ }
+ }
+ }
+
private static class TableCatalogResult {
final Table table;
diff --git a/core/src/main/java/org/apache/gravitino/meta/ColumnEntity.java
b/core/src/main/java/org/apache/gravitino/meta/ColumnEntity.java
index 5e68e48744..b46f07b983 100644
--- a/core/src/main/java/org/apache/gravitino/meta/ColumnEntity.java
+++ b/core/src/main/java/org/apache/gravitino/meta/ColumnEntity.java
@@ -27,6 +27,7 @@ import lombok.ToString;
import org.apache.gravitino.Audit;
import org.apache.gravitino.Auditable;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityFieldLimits;
import org.apache.gravitino.Field;
import org.apache.gravitino.rel.Column;
import org.apache.gravitino.rel.expressions.Expression;
@@ -40,12 +41,14 @@ import org.apache.gravitino.rel.types.Type;
public class ColumnEntity implements Entity, Auditable {
public static final Field ID = Field.required("id", Long.class, "The
column's unique identifier");
- public static final Field NAME = Field.required("name", String.class, "The
column's name");
+ public static final Field NAME =
+ Field.required("name", "The column's name",
EntityFieldLimits.MAX_NAME_LENGTH);
public static final Field POSITION =
Field.required("position", Integer.class, "The column's position");
public static final Field TYPE = Field.required("dataType", Type.class, "The
column's data type");
public static final Field COMMENT =
- Field.optional("comment", String.class, "The column's comment");
+ Field.optional(
+ "comment", "The column's comment",
EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH);
public static final Field NULLABLE =
Field.required("nullable", Boolean.class, "The column's nullable
property");
public static final Field AUTO_INCREMENT =
diff --git a/core/src/main/java/org/apache/gravitino/meta/TableEntity.java
b/core/src/main/java/org/apache/gravitino/meta/TableEntity.java
index 3995eec278..d95cd9a6a7 100644
--- a/core/src/main/java/org/apache/gravitino/meta/TableEntity.java
+++ b/core/src/main/java/org/apache/gravitino/meta/TableEntity.java
@@ -29,6 +29,7 @@ import lombok.ToString;
import lombok.experimental.Accessors;
import org.apache.gravitino.Auditable;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityFieldLimits;
import org.apache.gravitino.Field;
import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.Namespace;
@@ -46,7 +47,8 @@ import org.apache.gravitino.utils.CollectionUtils;
public class TableEntity implements Entity, Auditable, HasIdentifier {
public static final Field ID = Field.required("id", Long.class, "The table's
unique identifier");
- public static final Field NAME = Field.required("name", String.class, "The
table's name");
+ public static final Field NAME =
+ Field.required("name", "The table's name",
EntityFieldLimits.MAX_NAME_LENGTH);
public static final Field AUDIT_INFO =
Field.required("audit_info", AuditInfo.class, "The audit details of the
table");
public static final Field COLUMNS =
diff --git
a/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
b/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
index a84f6e152e..6df3796d16 100644
---
a/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
+++
b/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
@@ -52,6 +52,7 @@ import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.gravitino.Config;
import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityAlreadyExistsException;
+import org.apache.gravitino.EntityFieldLimits;
import org.apache.gravitino.GravitinoEnv;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
@@ -598,6 +599,179 @@ public class TestTableOperationDispatcher extends
TestOperationDispatcher {
reset(entityStore);
}
+ @Test
+ public void testRejectsOversizedTableNameBeforeExternalChange() throws
IOException {
+ Namespace tableNs = Namespace.of(metalake, catalog,
"schema_table_name_limit");
+ NameIdentifier validTableIdent = NameIdentifier.of(tableNs, "valid_table");
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ NameIdentifier oversizedTableIdent = NameIdentifier.of(tableNs,
oversizedName);
+ Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+ Column[] columns =
+ new Column[] {
+ TestColumn.builder()
+ .withName("col1")
+ .withPosition(0)
+ .withType(Types.StringType.get())
+ .build()
+ };
+
+
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()),
"comment", props);
+
+ IllegalArgumentException createException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.createTable(
+ oversizedTableIdent, columns, "comment", props, new
Transform[0]));
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
createException.getMessage());
+
+ tableOperationDispatcher.createTable(
+ validTableIdent, columns, "comment", props, new Transform[0]);
+ IllegalArgumentException renameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.alterTable(
+ validTableIdent, TableChange.rename(oversizedName)));
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
renameException.getMessage());
+
+ catalogManager.doWithCatalog(
+ NameIdentifier.of(metalake, catalog),
+ liveCatalog -> {
+ TestCatalogOperations testCatalogOperations =
(TestCatalogOperations) liveCatalog.ops();
+ Assertions.assertDoesNotThrow(() ->
testCatalogOperations.loadTable(validTableIdent));
+ Assertions.assertThrows(
+ NoSuchTableException.class,
+ () -> testCatalogOperations.loadTable(oversizedTableIdent));
+ return null;
+ });
+ }
+
+ @Test
+ public void testRejectsOversizedColumnFieldsBeforeExternalChange() throws
IOException {
+ Namespace tableNs = Namespace.of(metalake, catalog,
"schema_column_field_limits");
+ NameIdentifier validTableIdent = NameIdentifier.of(tableNs, "valid_table");
+ NameIdentifier invalidNameTableIdent = NameIdentifier.of(tableNs,
"invalid_column_name");
+ NameIdentifier invalidCommentTableIdent = NameIdentifier.of(tableNs,
"invalid_column_comment");
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ String oversizedComment =
"a".repeat(EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH + 1);
+ Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+ Column validColumn =
+ TestColumn.builder()
+ .withName("col1")
+ .withPosition(0)
+ .withType(Types.StringType.get())
+ .withComment("comment")
+ .build();
+
+
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()),
"comment", props);
+
+ IllegalArgumentException createNameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.createTable(
+ invalidNameTableIdent,
+ new Column[] {
+ TestColumn.builder()
+ .withName(oversizedName)
+ .withPosition(0)
+ .withType(Types.StringType.get())
+ .build()
+ },
+ "comment",
+ props,
+ new Transform[0]));
+ Assertions.assertEquals(
+ "The name of the column must not exceed 128 characters",
createNameException.getMessage());
+
+ IllegalArgumentException createCommentException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.createTable(
+ invalidCommentTableIdent,
+ new Column[] {
+ TestColumn.builder()
+ .withName("col1")
+ .withPosition(0)
+ .withType(Types.StringType.get())
+ .withComment(oversizedComment)
+ .build()
+ },
+ "comment",
+ props,
+ new Transform[0]));
+ Assertions.assertEquals(
+ "The comment of the column must not exceed 4096 characters",
+ createCommentException.getMessage());
+
+ tableOperationDispatcher.createTable(
+ validTableIdent, new Column[] {validColumn}, "comment", props, new
Transform[0]);
+
+ IllegalArgumentException addNameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.alterTable(
+ validTableIdent,
+ TableChange.addColumn(new String[] {oversizedName},
Types.StringType.get())));
+ Assertions.assertEquals(
+ "The name of the column must not exceed 128 characters",
addNameException.getMessage());
+
+ IllegalArgumentException addCommentException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.alterTable(
+ validTableIdent,
+ TableChange.addColumn(
+ new String[] {"col2"}, Types.StringType.get(),
oversizedComment)));
+ Assertions.assertEquals(
+ "The comment of the column must not exceed 4096 characters",
+ addCommentException.getMessage());
+
+ IllegalArgumentException renameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.alterTable(
+ validTableIdent,
+ TableChange.renameColumn(new String[] {"col1"},
oversizedName)));
+ Assertions.assertEquals(
+ "The name of the column must not exceed 128 characters",
renameException.getMessage());
+
+ IllegalArgumentException updateCommentException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tableOperationDispatcher.alterTable(
+ validTableIdent,
+ TableChange.updateColumnComment(new String[] {"col1"},
oversizedComment)));
+ Assertions.assertEquals(
+ "The comment of the column must not exceed 4096 characters",
+ updateCommentException.getMessage());
+
+ catalogManager.doWithCatalog(
+ NameIdentifier.of(metalake, catalog),
+ liveCatalog -> {
+ TestCatalogOperations testCatalogOperations =
(TestCatalogOperations) liveCatalog.ops();
+ Assertions.assertThrows(
+ NoSuchTableException.class,
+ () -> testCatalogOperations.loadTable(invalidNameTableIdent));
+ Assertions.assertThrows(
+ NoSuchTableException.class,
+ () -> testCatalogOperations.loadTable(invalidCommentTableIdent));
+ Column[] catalogColumns =
testCatalogOperations.loadTable(validTableIdent).columns();
+ Assertions.assertEquals(1, catalogColumns.length);
+ Assertions.assertEquals(validColumn.name(),
catalogColumns[0].name());
+ Assertions.assertEquals(validColumn.comment(),
catalogColumns[0].comment());
+ return null;
+ });
+ }
+
@Test
public void testCreateAndDropTable() throws IOException {
NameIdentifier tableIdent = NameIdentifier.of(metalake, catalog,
"schema71", "table31");
diff --git
a/core/src/test/java/org/apache/gravitino/meta/TestEntityFieldLimits.java
b/core/src/test/java/org/apache/gravitino/meta/TestEntityFieldLimits.java
index d2a77ae387..eab52d2474 100644
--- a/core/src/test/java/org/apache/gravitino/meta/TestEntityFieldLimits.java
+++ b/core/src/test/java/org/apache/gravitino/meta/TestEntityFieldLimits.java
@@ -39,6 +39,7 @@ import org.apache.gravitino.job.ShellJobTemplate;
import org.apache.gravitino.model.ModelVersion;
import org.apache.gravitino.policy.Policy;
import org.apache.gravitino.policy.PolicyContents;
+import org.apache.gravitino.rel.types.Types;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
@@ -75,6 +76,15 @@ public class TestEntityFieldLimits {
assertLengthLimit(builder, "alias", "model version",
EntityFieldLimits.MAX_NAME_LENGTH);
}
+ @Test
+ public void testColumnCommentLength() {
+ assertLengthLimit(
+ comment -> columnBuilder("column", comment),
+ "comment",
+ "column",
+ EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH);
+ }
+
@Test
public void testLengthCountsCodePoints() {
// An emoji is two UTF-16 chars but one character for MySQL (utf8mb4) and
PostgreSQL, where a
@@ -119,6 +129,17 @@ public class TestEntityFieldLimits {
private static Stream<Arguments> nameBuilders() {
return Stream.of(
+ Arguments.of("column", (Function<String, Entity>) name ->
columnBuilder(name, null)),
+ Arguments.of(
+ "table",
+ (Function<String, Entity>)
+ name ->
+ TableEntity.builder()
+ .withId(1L)
+ .withName(name)
+ .withNamespace(NAMESPACE)
+ .withAuditInfo(AuditInfo.EMPTY)
+ .build()),
Arguments.of("tag", (Function<String, Entity>) name ->
tagBuilder(name, null)),
Arguments.of(
"policy",
@@ -255,4 +276,15 @@ public class TestEntityFieldLimits {
.withAuditInfo(AuditInfo.EMPTY)
.build();
}
+
+ private static ColumnEntity columnBuilder(String name, String comment) {
+ return ColumnEntity.builder()
+ .withId(1L)
+ .withName(name)
+ .withPosition(0)
+ .withDataType(Types.StringType.get())
+ .withComment(comment)
+ .withAuditInfo(AuditInfo.EMPTY)
+ .build();
+ }
}
diff --git
a/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/ops/IcebergCatalogWrapper.java
b/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/ops/IcebergCatalogWrapper.java
index ed58b2885b..041fbca583 100644
---
a/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/ops/IcebergCatalogWrapper.java
+++
b/iceberg/iceberg-common/src/main/java/org/apache/gravitino/iceberg/common/ops/IcebergCatalogWrapper.java
@@ -36,13 +36,16 @@ import org.apache.gravitino.utils.ClassUtils;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.security.UserGroupInformation;
import org.apache.iceberg.BaseTable;
+import org.apache.iceberg.CatalogUtil;
import org.apache.iceberg.TableMetadata;
+import org.apache.iceberg.TableMetadataParser;
import org.apache.iceberg.Transaction;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.catalog.SupportsNamespaces;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.catalog.ViewCatalog;
+import org.apache.iceberg.io.FileIO;
import org.apache.iceberg.io.ResolvingFileIO;
import org.apache.iceberg.jdbc.JdbcCatalogWithMetadataLocationSupport;
import org.apache.iceberg.rest.CatalogHandlers;
@@ -248,6 +251,18 @@ public class IcebergCatalogWrapper implements
AutoCloseable {
return ((BaseTable)
getCatalog().loadTable(tableIdentifier)).operations().current();
}
+ /**
+ * Loads table metadata directly from a metadata file location.
+ *
+ * @param metadataLocation metadata file location
+ * @return parsed table metadata
+ */
+ public TableMetadata loadTableMetadataFromLocation(String metadataLocation) {
+ try (FileIO fileIO = CatalogUtil.loadFileIO(fileIOImpl(),
fileIOProperties(), null)) {
+ return TableMetadataParser.read(fileIO, metadataLocation);
+ }
+ }
+
/**
* Returns the FileIO implementation configured for this catalog.
*
diff --git
a/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/ops/TestIcebergCatalogWrapper.java
b/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/ops/TestIcebergCatalogWrapper.java
index e6928e56a8..48941d7f8f 100644
---
a/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/ops/TestIcebergCatalogWrapper.java
+++
b/iceberg/iceberg-common/src/test/java/org/apache/gravitino/iceberg/common/ops/TestIcebergCatalogWrapper.java
@@ -20,6 +20,8 @@ package org.apache.gravitino.iceberg.common.ops;
import java.io.IOException;
import java.lang.reflect.Method;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashMap;
import java.util.Map;
@@ -31,11 +33,15 @@ import org.apache.gravitino.iceberg.common.IcebergConfig;
import org.apache.gravitino.iceberg.common.cache.SupportsMetadataLocation;
import org.apache.gravitino.iceberg.common.cache.TableMetadataCache;
import org.apache.gravitino.iceberg.common.utils.IcebergCatalogUtil;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
import org.apache.iceberg.TableMetadata;
+import org.apache.iceberg.TableMetadataParser;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.rest.requests.CreateNamespaceRequest;
+import org.apache.iceberg.types.Types;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -104,6 +110,26 @@ public class TestIcebergCatalogWrapper {
Assertions.assertTrue(TrackingTableMetadataCache.CLOSED.get());
}
+ @Test
+ public void testLoadTableMetadataFromLocation(@TempDir Path tempDir) throws
Exception {
+ Schema schema = new Schema(Types.NestedField.required(1, "id",
Types.LongType.get()));
+ TableMetadata expected =
+ TableMetadata.newTableMetadata(
+ schema,
+ PartitionSpec.unpartitioned(),
+ tempDir.resolve("table").toUri().toString(),
+ Map.of());
+ Path metadataFile = tempDir.resolve("v1.metadata.json");
+ Files.writeString(metadataFile, TableMetadataParser.toJson(expected),
StandardCharsets.UTF_8);
+ IcebergCatalogWrapper wrapper =
+ new IcebergCatalogWrapper(
+ new IcebergConfig(Map.of(IcebergConstants.CATALOG_BACKEND,
"memory")));
+
+ TableMetadata actual =
wrapper.loadTableMetadataFromLocation(metadataFile.toUri().toString());
+
+ Assertions.assertEquals(schema.asStruct(), actual.schema().asStruct());
+ }
+
private static TableMetadataCache
invokeGetMetadataCache(IcebergCatalogWrapper wrapper)
throws Exception {
Method method =
IcebergCatalogWrapper.class.getDeclaredMethod("getMetadataCache");
diff --git
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergColumnFieldValidator.java
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergColumnFieldValidator.java
new file mode 100644
index 0000000000..5c4d78a4ef
--- /dev/null
+++
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergColumnFieldValidator.java
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino.iceberg.service.dispatcher;
+
+import static org.apache.gravitino.Entity.EntityType.COLUMN;
+
+import org.apache.gravitino.meta.ColumnEntity;
+import org.apache.iceberg.MetadataUpdate;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.rest.requests.UpdateTableRequest;
+
+final class IcebergColumnFieldValidator {
+
+ static void validateSchema(Schema schema) {
+ schema
+ .columns()
+ .forEach(
+ column -> {
+ ColumnEntity.NAME.validate(column.name(), COLUMN);
+ ColumnEntity.COMMENT.validate(column.doc(), COLUMN);
+ });
+ }
+
+ static void validateUpdate(UpdateTableRequest request) {
+ request.updates().stream()
+ .filter(MetadataUpdate.AddSchema.class::isInstance)
+ .map(MetadataUpdate.AddSchema.class::cast)
+ .map(MetadataUpdate.AddSchema::schema)
+ .forEach(IcebergColumnFieldValidator::validateSchema);
+ }
+
+ private IcebergColumnFieldValidator() {}
+}
diff --git
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergNamespaceOperationExecutor.java
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergNamespaceOperationExecutor.java
index cdbc9a24ca..a884685c92 100644
---
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergNamespaceOperationExecutor.java
+++
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergNamespaceOperationExecutor.java
@@ -22,11 +22,15 @@ package org.apache.gravitino.iceberg.service.dispatcher;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
+import org.apache.gravitino.Entity;
import org.apache.gravitino.auth.AuthConstants;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
+import org.apache.gravitino.iceberg.service.CatalogWrapperForREST;
import org.apache.gravitino.iceberg.service.IcebergCatalogWrapperManager;
import org.apache.gravitino.iceberg.service.cleanup.IcebergCleanupManager;
import org.apache.gravitino.listener.api.event.IcebergRequestContext;
+import org.apache.gravitino.meta.TableEntity;
+import org.apache.iceberg.TableMetadata;
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.rest.requests.CreateNamespaceRequest;
import org.apache.iceberg.rest.requests.RegisterTableRequest;
@@ -128,12 +132,17 @@ public class IcebergNamespaceOperationExecutor implements
IcebergNamespaceOperat
IcebergRequestContext context,
Namespace namespace,
RegisterTableRequest registerTableRequest) {
+ TableEntity.NAME.validate(registerTableRequest.name(),
Entity.EntityType.TABLE);
IcebergCleanupHelper.rejectIfBeingPurged(
cleanupManager, context.catalogName(), namespace,
registerTableRequest.name());
- return icebergCatalogWrapperManager
- .getCatalogWrapper(context.catalogName())
- .registerTable(namespace, registerTableRequest,
context.requestCredentialVending());
+ CatalogWrapperForREST catalogWrapper =
+ icebergCatalogWrapperManager.getCatalogWrapper(context.catalogName());
+ TableMetadata metadata =
+
catalogWrapper.loadTableMetadataFromLocation(registerTableRequest.metadataLocation());
+ IcebergColumnFieldValidator.validateSchema(metadata.schema());
+ return catalogWrapper.registerTable(
+ namespace, registerTableRequest, context.requestCredentialVending());
}
@Override
diff --git
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergTableOperationExecutor.java
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergTableOperationExecutor.java
index 5a8a21d71e..c935f6d700 100644
---
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergTableOperationExecutor.java
+++
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/dispatcher/IcebergTableOperationExecutor.java
@@ -34,10 +34,12 @@ import
org.apache.gravitino.iceberg.service.authorization.IcebergRESTServerConte
import org.apache.gravitino.iceberg.service.cleanup.IcebergCleanupJob;
import org.apache.gravitino.iceberg.service.cleanup.IcebergCleanupManager;
import org.apache.gravitino.listener.api.event.IcebergRequestContext;
+import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.server.authorization.MetadataAuthzHelper;
import
org.apache.gravitino.server.authorization.expression.AuthorizationExpressionConstants;
import org.apache.gravitino.utils.HierarchicalSchemaUtil;
import org.apache.iceberg.TableMetadata;
+import org.apache.iceberg.UpdateRequirement;
import org.apache.iceberg.catalog.Namespace;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.rest.requests.CreateTableRequest;
@@ -70,6 +72,8 @@ public class IcebergTableOperationExecutor implements
IcebergTableOperationDispa
@Override
public LoadTableResponse createTable(
IcebergRequestContext context, Namespace namespace, CreateTableRequest
createTableRequest) {
+ TableEntity.NAME.validate(createTableRequest.name(),
Entity.EntityType.TABLE);
+ IcebergColumnFieldValidator.validateSchema(createTableRequest.schema());
IcebergCleanupHelper.rejectIfBeingPurged(
cleanupManager, context.catalogName(), namespace,
createTableRequest.name());
@@ -115,6 +119,11 @@ public class IcebergTableOperationExecutor implements
IcebergTableOperationDispa
IcebergRequestContext context,
TableIdentifier tableIdentifier,
UpdateTableRequest updateTableRequest) {
+ if (updateTableRequest.requirements().stream()
+
.anyMatch(UpdateRequirement.AssertTableDoesNotExist.class::isInstance)) {
+ TableEntity.NAME.validate(tableIdentifier.name(),
Entity.EntityType.TABLE);
+ }
+ IcebergColumnFieldValidator.validateUpdate(updateTableRequest);
return icebergCatalogWrapperManager
.getCatalogWrapper(context.catalogName())
.updateTable(tableIdentifier, updateTableRequest);
@@ -189,6 +198,7 @@ public class IcebergTableOperationExecutor implements
IcebergTableOperationDispa
@Override
public void renameTable(IcebergRequestContext context, RenameTableRequest
renameTableRequest) {
+ TableEntity.NAME.validate(renameTableRequest.destination().name(),
Entity.EntityType.TABLE);
icebergCatalogWrapperManager
.getCatalogWrapper(context.catalogName())
.renameTable(renameTableRequest);
diff --git
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergAsyncPurge.java
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergAsyncPurge.java
index 7b39b9bd71..02ef946745 100644
---
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergAsyncPurge.java
+++
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergAsyncPurge.java
@@ -197,11 +197,15 @@ class TestIcebergAsyncPurge {
CatalogWrapperForREST wrapper = mock(CatalogWrapperForREST.class);
IcebergCleanupManager cleanup = mock(IcebergCleanupManager.class);
when(cleanup.isNameOccupied(CATALOG_ID, "db", "t")).thenReturn(false);
+ TableMetadata metadata = mock(TableMetadata.class);
+ when(metadata.schema()).thenReturn(SCHEMA);
+
when(wrapper.loadTableMetadataFromLocation("s3://b/db/t/metadata/0.json")).thenReturn(metadata);
try (MockedStatic<GravitinoEnv> ignored = mockCatalogId()) {
namespaceExecutor(wrapper, Optional.of(cleanup))
.registerTable(context(false), DB, registerReq());
}
+
verify(wrapper).loadTableMetadataFromLocation("s3://b/db/t/metadata/0.json");
verify(wrapper).registerTable(any(), any(), anyBoolean());
}
diff --git
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergNamespaceOperationExecutor.java
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergNamespaceOperationExecutor.java
index d1a45f0ab2..0b37b8e3fc 100644
---
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergNamespaceOperationExecutor.java
+++
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergNamespaceOperationExecutor.java
@@ -21,28 +21,45 @@ package org.apache.gravitino.iceberg.service.dispatcher;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
+import org.apache.gravitino.EntityFieldLimits;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
+import org.apache.gravitino.iceberg.common.IcebergConfig;
import org.apache.gravitino.iceberg.service.CatalogWrapperForREST;
import org.apache.gravitino.iceberg.service.IcebergCatalogWrapperManager;
import org.apache.gravitino.listener.api.event.IcebergRequestContext;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.TableMetadata;
+import org.apache.iceberg.TableMetadataParser;
import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.rest.requests.CreateNamespaceRequest;
+import org.apache.iceberg.rest.requests.ImmutableRegisterTableRequest;
+import org.apache.iceberg.rest.requests.RegisterTableRequest;
import org.apache.iceberg.rest.requests.RegisterViewRequest;
import org.apache.iceberg.rest.responses.CreateNamespaceResponse;
import org.apache.iceberg.rest.responses.GetNamespaceResponse;
import org.apache.iceberg.rest.responses.ListNamespacesResponse;
import org.apache.iceberg.rest.responses.LoadViewResponse;
+import org.apache.iceberg.types.Types.NestedField;
+import org.apache.iceberg.types.Types.StringType;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
import org.mockito.ArgumentCaptor;
public class TestIcebergNamespaceOperationExecutor {
@@ -191,6 +208,98 @@ public class TestIcebergNamespaceOperationExecutor {
Assertions.assertEquals(mockResponse, result);
}
+ @Test
+ public void testRejectsOversizedTableNameBeforeRegister() {
+ RegisterTableRequest request = mock(RegisterTableRequest.class);
+
when(request.name()).thenReturn("a".repeat(EntityFieldLimits.MAX_NAME_LENGTH +
1));
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> executor.registerTable(mockContext,
Namespace.of("test_namespace"), request));
+
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
exception.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ @Test
+ public void testRejectsOversizedColumnNameBeforeRegister() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ Schema schema = new Schema(NestedField.required(1, oversizedName,
StringType.get()));
+ assertRegisterRejectsSchema(schema, "The name of the column must not
exceed 128 characters");
+ }
+
+ @Test
+ public void testRejectsOversizedColumnCommentBeforeRegister() {
+ String oversizedComment =
"a".repeat(EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH + 1);
+ Schema schema = new Schema(NestedField.required(1, "col1",
StringType.get(), oversizedComment));
+ assertRegisterRejectsSchema(
+ schema, "The comment of the column must not exceed 4096 characters");
+ }
+
+ @Test
+ public void testInvalidRegisterOverwriteLeavesMetadataUnchanged(@TempDir
Path tempDir)
+ throws Exception {
+ IcebergConfig config =
+ new IcebergConfig(
+ Map.of(
+ IcebergConstants.CATALOG_BACKEND,
+ "jdbc",
+ IcebergConstants.URI,
+ "jdbc:sqlite:" + tempDir.resolve("catalog.db"),
+ IcebergConstants.WAREHOUSE,
+ tempDir.resolve("warehouse").toString(),
+ IcebergConstants.GRAVITINO_JDBC_DRIVER,
+ "org.sqlite.JDBC",
+ IcebergConstants.ICEBERG_JDBC_USER,
+ "test",
+ IcebergConstants.ICEBERG_JDBC_PASSWORD,
+ "test",
+ IcebergConstants.ICEBERG_JDBC_INITIALIZE,
+ "true"));
+ CatalogWrapperForREST catalogWrapper = new
CatalogWrapperForREST("test_catalog", config);
+ try {
+
when(mockWrapperManager.getCatalogWrapper("test_catalog")).thenReturn(catalogWrapper);
+ Namespace namespace = Namespace.of("test_namespace");
+ catalogWrapper.createNamespace(
+ CreateNamespaceRequest.builder().withNamespace(namespace).build());
+ String originalMetadataLocation =
+ writeMetadata(
+ tempDir.resolve("v1.metadata.json"),
+ new Schema(NestedField.required(1, "id", StringType.get())));
+ RegisterTableRequest originalRequest =
+ ImmutableRegisterTableRequest.builder()
+ .name("test_table")
+ .metadataLocation(originalMetadataLocation)
+ .build();
+ catalogWrapper.registerTable(namespace, originalRequest, false);
+
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ String invalidMetadataLocation =
+ writeMetadata(
+ tempDir.resolve("v2.metadata.json"),
+ new Schema(NestedField.required(1, oversizedName,
StringType.get())));
+ RegisterTableRequest invalidOverwriteRequest =
+ ImmutableRegisterTableRequest.builder()
+ .name("test_table")
+ .metadataLocation(invalidMetadataLocation)
+ .overwrite(true)
+ .build();
+
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> executor.registerTable(mockContext, namespace,
invalidOverwriteRequest));
+
+ Assertions.assertEquals(
+ Optional.of(originalMetadataLocation),
+ catalogWrapper.getTableMetadataLocation(
+ TableIdentifier.of(namespace, originalRequest.name())));
+ } finally {
+ catalogWrapper.close();
+ }
+ }
+
@Test
public void testDropNestedNamespacePassesCorrectLevels() {
Namespace nestedNs = Namespace.of("A", "B", "C");
@@ -260,4 +369,36 @@ public class TestIcebergNamespaceOperationExecutor {
verify(mockCatalogWrapper).namespaceExists(ns);
Assertions.assertFalse(exists);
}
+
+ private void assertRegisterRejectsSchema(Schema schema, String
expectedMessage) {
+ Namespace namespace = Namespace.of("test_namespace");
+ String metadataLocation = "file:/tmp/test.metadata.json";
+ RegisterTableRequest request = mock(RegisterTableRequest.class);
+ when(request.name()).thenReturn("test_table");
+ when(request.metadataLocation()).thenReturn(metadataLocation);
+ TableMetadata metadata =
+ TableMetadata.newTableMetadata(
+ schema, PartitionSpec.unpartitioned(), "file:/tmp/table",
Collections.emptyMap());
+
when(mockCatalogWrapper.loadTableMetadataFromLocation(metadataLocation)).thenReturn(metadata);
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> executor.registerTable(mockContext, namespace, request));
+
+ Assertions.assertEquals(expectedMessage, exception.getMessage());
+ verify(mockCatalogWrapper).loadTableMetadataFromLocation(metadataLocation);
+ verify(mockCatalogWrapper, never()).registerTable(namespace, request,
false);
+ }
+
+ private static String writeMetadata(Path metadataFile, Schema schema) throws
Exception {
+ TableMetadata metadata =
+ TableMetadata.newTableMetadata(
+ schema,
+ PartitionSpec.unpartitioned(),
+ metadataFile.getParent().resolve("table").toUri().toString(),
+ Collections.emptyMap());
+ Files.writeString(metadataFile, TableMetadataParser.toJson(metadata),
StandardCharsets.UTF_8);
+ return metadataFile.toUri().toString();
+ }
}
diff --git
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergTableOperationExecutor.java
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergTableOperationExecutor.java
index 9c5ebe158a..1e7f0f2a70 100644
---
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergTableOperationExecutor.java
+++
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/dispatcher/TestIcebergTableOperationExecutor.java
@@ -23,18 +23,26 @@ import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
+import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
+import org.apache.gravitino.EntityFieldLimits;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
import org.apache.gravitino.iceberg.service.CatalogWrapperForREST;
import org.apache.gravitino.iceberg.service.IcebergCatalogWrapperManager;
import org.apache.gravitino.listener.api.event.IcebergRequestContext;
+import org.apache.iceberg.MetadataUpdate;
import org.apache.iceberg.Schema;
+import org.apache.iceberg.UpdateRequirement;
import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.rest.requests.CreateTableRequest;
+import org.apache.iceberg.rest.requests.RenameTableRequest;
+import org.apache.iceberg.rest.requests.UpdateTableRequest;
import org.apache.iceberg.rest.responses.LoadTableResponse;
import org.apache.iceberg.types.Types.NestedField;
import org.apache.iceberg.types.Types.StringType;
@@ -189,4 +197,135 @@ public class TestIcebergTableOperationExecutor {
requestCaptor.getValue().stageCreate(),
"stageCreate=false must remain false when rebuilding request");
}
+
+ @Test
+ public void testRejectsOversizedTableNameBeforeCreate() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ CreateTableRequest request =
+
CreateTableRequest.builder().withName(oversizedName).withSchema(TABLE_SCHEMA).build();
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> executor.createTable(mockContext,
Namespace.of("test_namespace"), request));
+
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
exception.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ @Test
+ public void testRejectsOversizedTableNameBeforeRename() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ RenameTableRequest request =
+ RenameTableRequest.builder()
+ .withSource(TableIdentifier.of("test_namespace", "source"))
+ .withDestination(TableIdentifier.of("test_namespace",
oversizedName))
+ .build();
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class, () ->
executor.renameTable(mockContext, request));
+
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
exception.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ @Test
+ public void testRejectsOversizedColumnFieldsBeforeCreate() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ Schema oversizedNameSchema =
+ new Schema(NestedField.required(1, oversizedName, StringType.get()));
+ CreateTableRequest oversizedNameRequest =
+
CreateTableRequest.builder().withName("test_table").withSchema(oversizedNameSchema).build();
+
+ IllegalArgumentException nameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ executor.createTable(
+ mockContext, Namespace.of("test_namespace"),
oversizedNameRequest));
+ Assertions.assertEquals(
+ "The name of the column must not exceed 128 characters",
nameException.getMessage());
+
+ String oversizedComment =
"a".repeat(EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH + 1);
+ Schema oversizedCommentSchema =
+ new Schema(NestedField.required(1, "col1", StringType.get(),
oversizedComment));
+ CreateTableRequest oversizedCommentRequest =
+ CreateTableRequest.builder()
+ .withName("test_table")
+ .withSchema(oversizedCommentSchema)
+ .build();
+
+ IllegalArgumentException commentException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ executor.createTable(
+ mockContext, Namespace.of("test_namespace"),
oversizedCommentRequest));
+ Assertions.assertEquals(
+ "The comment of the column must not exceed 4096 characters",
commentException.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ @Test
+ public void testRejectsOversizedColumnFieldsBeforeUpdate() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ UpdateTableRequest oversizedNameRequest =
+ updateRequest(new Schema(NestedField.required(1, oversizedName,
StringType.get())));
+
+ IllegalArgumentException nameException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ executor.updateTable(
+ mockContext,
+ TableIdentifier.of("test_namespace", "test_table"),
+ oversizedNameRequest));
+ Assertions.assertEquals(
+ "The name of the column must not exceed 128 characters",
nameException.getMessage());
+
+ String oversizedComment =
"a".repeat(EntityFieldLimits.MAX_COLUMN_COMMENT_LENGTH + 1);
+ UpdateTableRequest oversizedCommentRequest =
+ updateRequest(
+ new Schema(NestedField.required(1, "col1", StringType.get(),
oversizedComment)));
+
+ IllegalArgumentException commentException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ executor.updateTable(
+ mockContext,
+ TableIdentifier.of("test_namespace", "test_table"),
+ oversizedCommentRequest));
+ Assertions.assertEquals(
+ "The comment of the column must not exceed 4096 characters",
commentException.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ @Test
+ public void testRejectsOversizedTableNameBeforeStagedCreateCommit() {
+ String oversizedName = "a".repeat(EntityFieldLimits.MAX_NAME_LENGTH + 1);
+ UpdateTableRequest request =
+ new UpdateTableRequest(
+ Collections.singletonList(new
UpdateRequirement.AssertTableDoesNotExist()),
+ Collections.emptyList());
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ executor.updateTable(
+ mockContext, TableIdentifier.of("test_namespace",
oversizedName), request));
+
+ Assertions.assertEquals(
+ "The name of the table must not exceed 128 characters",
exception.getMessage());
+ verifyNoInteractions(mockCatalogWrapper);
+ }
+
+ private static UpdateTableRequest updateRequest(Schema schema) {
+ return new UpdateTableRequest(
+ Collections.emptyList(), Collections.singletonList(new
MetadataUpdate.AddSchema(schema)));
+ }
}
diff --git
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/CatalogWrapperForTest.java
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/CatalogWrapperForTest.java
index 9786d7b2fc..9affd68e75 100644
---
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/CatalogWrapperForTest.java
+++
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/CatalogWrapperForTest.java
@@ -82,15 +82,7 @@ public class CatalogWrapperForTest extends
CatalogWrapperForREST {
// metadata.json file at the given location), so build a mock
LoadTableResponse here.
// Honor cloud URIs (e.g. s3://) in metadataLocation so credential vending
tests can
// verify the vended path; default to /mock otherwise for existing tests.
- String location =
- request.metadataLocation().contains("://") ?
request.metadataLocation() : "/mock";
- Schema mockSchema = new Schema(NestedField.of(1, false, "foo_string",
StringType.get()));
- TableMetadata baseMetadata =
- TableMetadata.newTableMetadata(
- mockSchema, PartitionSpec.unpartitioned(), location,
ImmutableMap.of());
- String json = TableMetadataParser.toJson(baseMetadata);
- TableMetadata tableMetadata =
- TableMetadataParser.fromJson(location + "/metadata/v1.metadata.json",
json);
+ TableMetadata tableMetadata =
mockTableMetadata(request.metadataLocation());
LoadTableResponse loadTableResponse =
LoadTableResponse.builder()
.withTableMetadata(tableMetadata)
@@ -113,6 +105,11 @@ public class CatalogWrapperForTest extends
CatalogWrapperForREST {
return loadTableResponse;
}
+ @Override
+ public TableMetadata loadTableMetadataFromLocation(String metadataLocation) {
+ return mockTableMetadata(metadataLocation);
+ }
+
@Override
public LoadViewResponse registerView(Namespace namespace,
RegisterViewRequest request) {
if (request.name().contains("fail")) {
@@ -157,6 +154,16 @@ public class CatalogWrapperForTest extends
CatalogWrapperForREST {
request.properties().getOrDefault(GENERATE_PLAN_TASKS_DATA_PROP,
Boolean.FALSE.toString()));
}
+ private static TableMetadata mockTableMetadata(String metadataLocation) {
+ String location = metadataLocation.contains("://") ? metadataLocation :
"/mock";
+ Schema mockSchema = new Schema(NestedField.of(1, false, "foo_string",
StringType.get()));
+ TableMetadata baseMetadata =
+ TableMetadata.newTableMetadata(
+ mockSchema, PartitionSpec.unpartitioned(), location,
ImmutableMap.of());
+ String json = TableMetadataParser.toJson(baseMetadata);
+ return TableMetadataParser.fromJson(location +
"/metadata/v1.metadata.json", json);
+ }
+
private void appendSampleData(Namespace namespace, String tableName) {
try {
Table table = getCatalog().loadTable(TableIdentifier.of(namespace,
tableName));