This is an automated email from the ASF dual-hosted git repository.
roryqi pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 847266d0a4 [Cherry-pick to branch-1.3] [#13565] fix: Validate table
and column field lengths (#13551) (#13572)
847266d0a4 is described below
commit 847266d0a40cb9197938f4575592f7ef671a1f57
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Sep 29 13:10:05 2026 +0800
[Cherry-pick to branch-1.3] [#13565] fix: Validate table and column field
lengths (#13551) (#13572)
**Cherry-pick Information:**
- Original commit: 8e9ca0009d68dbf9edb5e313d3d51688b9c4cff1
- Target branch: branch-1.3
- Status: conflicts resolved
### What changes were proposed in this pull request?
- Validate table and top-level column names against the 128-character
relational-store limit.
- Validate top-level column comments against the branch-1.3
relational-store limit of 256 characters.
- Perform validation before external catalog mutations.
- Resolve the Iceberg REST test conflicts without backporting unrelated
view-registration code.
### Why are the changes needed?
Oversized values could mutate the external catalog before Gravitino
metadata persistence failed. The branch-1.3 H2, MySQL, and PostgreSQL
schemas use VARCHAR(256) for column comments, so this backport uses 256
rather than the main-branch limit.
Fix: #13565
### 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
- git diff --check
---------
Co-authored-by: roryqi <[email protected]>
Co-authored-by: roryqi <[email protected]>
---
.../catalog/TableOperationDispatcher.java | 33 ++++
.../org/apache/gravitino/meta/ColumnEntity.java | 6 +-
.../org/apache/gravitino/meta/TableEntity.java | 4 +-
.../catalog/TestTableOperationDispatcher.java | 175 +++++++++++++++++++++
.../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 | 140 +++++++++++++++++
.../TestIcebergTableOperationExecutor.java | 139 ++++++++++++++++
.../service/rest/CatalogWrapperForTest.java | 25 +--
14 files changed, 658 insertions(+), 15 deletions(-)
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 52297b0383..43cb76ff7c 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;
@@ -209,6 +210,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());
@@ -259,9 +263,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);
@@ -1011,6 +1017,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..a5738a0f10 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,13 @@ 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_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 8110f5a6ac..55b5678461 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 c576682425..0a819235bc 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;
@@ -61,6 +62,7 @@ import org.apache.gravitino.auth.AuthConstants;
import org.apache.gravitino.connector.TestCatalogOperations;
import org.apache.gravitino.dto.util.DTOConverters;
import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.exceptions.NoSuchTableException;
import org.apache.gravitino.lock.LockManager;
import org.apache.gravitino.lock.LockType;
import org.apache.gravitino.meta.AuditInfo;
@@ -503,6 +505,179 @@ public class TestTableOperationDispatcher extends
TestOperationDispatcher {
Assertions.assertEquals("test", alteredTable4.auditInfo().lastModifier());
}
+ @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_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 256 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 256 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 256 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..25edbed72e 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_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 566c5289d0..95ad00836a 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
@@ -39,13 +39,16 @@ import org.apache.gravitino.utils.IsolatedClassLoader;
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;
@@ -254,6 +257,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 364ba55120..d69b5a30e4 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;
@@ -126,11 +130,16 @@ 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());
}
}
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 6307013315..d95ee792e8 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;
@@ -68,6 +70,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());
@@ -113,6 +117,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);
@@ -187,6 +196,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 2e0e1b5a8d..7fa75a4df8 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,26 +21,43 @@ 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.responses.CreateNamespaceResponse;
import org.apache.iceberg.rest.responses.GetNamespaceResponse;
import org.apache.iceberg.rest.responses.ListNamespacesResponse;
+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 {
@@ -176,6 +193,97 @@ 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_COMMENT_LENGTH
+ 1);
+ Schema schema = new Schema(NestedField.required(1, "col1",
StringType.get(), oversizedComment));
+ assertRegisterRejectsSchema(schema, "The comment of the column must not
exceed 256 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");
@@ -245,4 +353,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..a8c7f7b230 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_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 256 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_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 256 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 5af252880f..ab50d014e0 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
@@ -74,15 +74,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)
@@ -105,6 +97,11 @@ public class CatalogWrapperForTest extends
CatalogWrapperForREST {
return loadTableResponse;
}
+ @Override
+ public TableMetadata loadTableMetadataFromLocation(String metadataLocation) {
+ return mockTableMetadata(metadataLocation);
+ }
+
private boolean shouldGeneratePlanTasksData(CreateTableRequest request) {
if (request.properties() == null) {
return false;
@@ -113,6 +110,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));