From 767658c920c57c45e0ba55d7ba420de19ad259c2 Mon Sep 17 00:00:00 2001 From: Levi Jiang Date: Wed, 24 Jun 2026 13:18:00 -0700 Subject: [PATCH 1/2] feat(tables): gate RTAS per-table behind 'replace.enabled' table property Add a per-table flag that gates the RTAS (Replace Table As Select) command. RTAS is now blocked by default and is only allowed when the table has the regular table property 'replace.enabled' set to true. The gate is enforced in TablesServiceImpl (stageReplace) and IcebergSnapshotsServiceImpl (replaceCommit) via AuthorizationUtils.checkReplaceTableEnabled, which reads the table's persisted properties (so it cannot be bypassed by the incoming request body) and returns HTTP 400 when RTAS is disabled. Updated existing replace tests to opt in, added a negative test, and opted in the Spark/app RTAS integration tests. Co-Authored-By: Claude Opus 4.8 --- .../openhouse/catalog/e2e/RTASJavaTest.java | 6 +++ .../openhouse/jobs/spark/OperationsTest.java | 6 +++ .../openhouse/spark/catalogtest/RTASTest.java | 16 ++++++++ .../UnsupportedClientOperationException.java | 3 +- .../services/IcebergSnapshotsServiceImpl.java | 2 + .../tables/services/TablesServiceImpl.java | 2 + .../tables/utils/AuthorizationUtils.java | 29 +++++++++++++ .../e2e/h2/SnapshotsControllerTest.java | 7 ++++ .../tables/e2e/h2/TablesControllerTest.java | 41 ++++++++++++++++++- 9 files changed, 110 insertions(+), 2 deletions(-) diff --git a/apps/spark/src/test/java/com/linkedin/openhouse/catalog/e2e/RTASJavaTest.java b/apps/spark/src/test/java/com/linkedin/openhouse/catalog/e2e/RTASJavaTest.java index 23cd7ae8d..35040547c 100644 --- a/apps/spark/src/test/java/com/linkedin/openhouse/catalog/e2e/RTASJavaTest.java +++ b/apps/spark/src/test/java/com/linkedin/openhouse/catalog/e2e/RTASJavaTest.java @@ -65,6 +65,9 @@ void testReplaceTransaction() throws Exception { Table table = catalog.buildTable(TABLE_IDENT, SCHEMA).withPartitionSpec(SPEC).create(); table.newAppend().appendFile(FILE_A).commit(); + // RTAS is disabled by default; opt the table in before replacing it. + table.updateProperties().set("replace.enabled", "true").commit(); + String originalLocation = table.location(); String originalMetadataLocation = getMetadataLocation(table); long originalSnapshotId = table.currentSnapshot().snapshotId(); @@ -120,6 +123,9 @@ void testCreateOrReplaceTransaction() throws Exception { // verify that the table was created with one snapshot assertEquals(1, Iterables.size(table.snapshots()), "Should have one snapshot after create"); + // RTAS is disabled by default; opt the table in before replacing it. + table.updateProperties().set("replace.enabled", "true").commit(); + // create or replace on existing table should replace it Transaction replaceTxn = catalog diff --git a/apps/spark/src/test/java/com/linkedin/openhouse/jobs/spark/OperationsTest.java b/apps/spark/src/test/java/com/linkedin/openhouse/jobs/spark/OperationsTest.java index 2e44c9263..a32274662 100644 --- a/apps/spark/src/test/java/com/linkedin/openhouse/jobs/spark/OperationsTest.java +++ b/apps/spark/src/test/java/com/linkedin/openhouse/jobs/spark/OperationsTest.java @@ -752,6 +752,12 @@ public void testSnapshotsExpirationAfterReplaceTable() throws Exception { prepareTable(ops, tableName); populateTable(ops, tableName, numInserts); + // RTAS is disabled by default; opt the table in before replacing it. + ops.spark() + .sql( + String.format( + "ALTER TABLE %s SET TBLPROPERTIES ('replace.enabled'='true')", tableName)); + // replace the table using RTAS ops.spark() .sql( diff --git a/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java b/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java index e831ca706..ad921977e 100644 --- a/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java +++ b/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java @@ -60,6 +60,10 @@ public void testRTAS() throws Exception { spark.sql(String.format("ALTER TABLE %s SET POLICY (HISTORY MAX_AGE=24H)", tableName)); + // RTAS is disabled by default; opt the table in before replacing it. + spark.sql( + String.format("ALTER TABLE %s SET TBLPROPERTIES ('replace.enabled'='true')", tableName)); + String expectedTableLocation = catalog.loadTable(tableIdent).location(); // replace table @@ -125,6 +129,10 @@ public void testCreateRTAS() throws Exception { String expectedTableLocation = catalog.loadTable(tableIdent).location(); + // RTAS is disabled by default; opt the table in before replacing it. + spark.sql( + String.format("ALTER TABLE %s SET TBLPROPERTIES ('replace.enabled'='true')", tableName)); + // create or replace table should replace the table spark.sql( String.format( @@ -169,6 +177,10 @@ public void testDataFrameV2Replace() throws Exception { spark.table(sourceName).writeTo(tableName).using("iceberg").create(); + // RTAS is disabled by default; opt the table in before replacing it. + spark.sql( + String.format("ALTER TABLE %s SET TBLPROPERTIES ('replace.enabled'='true')", tableName)); + String expectedTableLocation = catalog.loadTable(tableIdent).location(); spark @@ -226,6 +238,10 @@ public void testDataFrameV2CreateOrReplace() throws Exception { .using("iceberg") .createOrReplace(); + // RTAS is disabled by default; opt the table in before replacing it. + spark.sql( + String.format("ALTER TABLE %s SET TBLPROPERTIES ('replace.enabled'='true')", tableName)); + String expectedTableLocation = catalog.loadTable(tableIdent).location(); spark diff --git a/services/common/src/main/java/com/linkedin/openhouse/common/exception/UnsupportedClientOperationException.java b/services/common/src/main/java/com/linkedin/openhouse/common/exception/UnsupportedClientOperationException.java index b76fbc8da..98ab0cba3 100644 --- a/services/common/src/main/java/com/linkedin/openhouse/common/exception/UnsupportedClientOperationException.java +++ b/services/common/src/main/java/com/linkedin/openhouse/common/exception/UnsupportedClientOperationException.java @@ -20,6 +20,7 @@ public enum Operation { GRANT_ON_UNSHARED_TABLES, ALTER_TABLE_TYPE, GRANT_ON_LOCKED_TABLES, - LOCKED_TABLE_OPERATION + LOCKED_TABLE_OPERATION, + RTAS_ON_DISABLED_TABLE } } diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/services/IcebergSnapshotsServiceImpl.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/services/IcebergSnapshotsServiceImpl.java index af43a169e..9d9fa0d63 100644 --- a/services/tables/src/main/java/com/linkedin/openhouse/tables/services/IcebergSnapshotsServiceImpl.java +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/services/IcebergSnapshotsServiceImpl.java @@ -67,6 +67,8 @@ public Pair putIcebergSnapshots( if (tableDto.isPresent() && icebergSnapshotRequestBody.getCreateUpdateTableRequestBody().isReplaceCommit()) { + // RTAS is gated per-table and disabled by default; ensure it is enabled on the table. + authorizationUtils.checkReplaceTableEnabled(tableDto.get()); // Check if table creator has the privilege to replace the table. authorizationUtils.checkReplaceTablePrivilege(tableDto.get(), tableCreatorUpdater); } else if (tableDto.isPresent()) { diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/services/TablesServiceImpl.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/services/TablesServiceImpl.java index 2c6bf611b..53c5a3481 100644 --- a/services/tables/src/main/java/com/linkedin/openhouse/tables/services/TablesServiceImpl.java +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/services/TablesServiceImpl.java @@ -100,6 +100,8 @@ public Pair putTable( // Special case handling if (tableDto.isPresent() && createUpdateTableRequestBody.isStageReplace()) { + // RTAS is gated per-table and disabled by default; ensure it is enabled on the table. + authorizationUtils.checkReplaceTableEnabled(tableDto.get()); // Check if table creator has the privilege to replace the table. authorizationUtils.checkReplaceTablePrivilege(tableDto.get(), tableCreatorUpdater); } else if (tableDto.isPresent()) { diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/utils/AuthorizationUtils.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/utils/AuthorizationUtils.java index 659c72a42..822db6179 100644 --- a/services/tables/src/main/java/com/linkedin/openhouse/tables/utils/AuthorizationUtils.java +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/utils/AuthorizationUtils.java @@ -1,10 +1,12 @@ package com.linkedin.openhouse.tables.utils; +import com.linkedin.openhouse.common.exception.UnsupportedClientOperationException; import com.linkedin.openhouse.tables.authorization.AuthorizationHandler; import com.linkedin.openhouse.tables.authorization.Privileges; import com.linkedin.openhouse.tables.common.TableType; import com.linkedin.openhouse.tables.model.DatabaseDto; import com.linkedin.openhouse.tables.model.TableDto; +import java.util.Map; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.security.access.AccessDeniedException; @@ -15,6 +17,12 @@ @Component public class AuthorizationUtils { + /** + * Table property that gates RTAS (Replace Table As Select). RTAS is disabled by default; a table + * can only be replaced when this property is explicitly set to {@code true}. + */ + public static final String RTAS_ENABLED_TABLE_PROP = "replace.enabled"; + @Autowired AuthorizationHandler authorizationHandler; /** @@ -118,4 +126,25 @@ public void checkReplaceTablePrivilege(TableDto tableDto, String actingPrincipal actingPrincipal)); } } + + /** + * Checks if RTAS (Replace Table As Select) is enabled for the table. RTAS is disabled by default + * and a table can only be replaced once its {@value #RTAS_ENABLED_TABLE_PROP} table property has + * been explicitly set to {@code true}. The check is performed against the table's persisted + * properties so the gate cannot be bypassed by the incoming request body. + * + * @param tableDto the existing table targeted by the replace operation + */ + public void checkReplaceTableEnabled(TableDto tableDto) { + Map tableProperties = tableDto.getTableProperties(); + if (tableProperties == null + || !Boolean.parseBoolean(tableProperties.get(RTAS_ENABLED_TABLE_PROP))) { + throw new UnsupportedClientOperationException( + UnsupportedClientOperationException.Operation.RTAS_ON_DISABLED_TABLE, + String.format( + "Replace (RTAS) is not enabled for table %s.%s. Set the table property '%s=true' " + + "before issuing a replace.", + tableDto.getDatabaseId(), tableDto.getTableId(), RTAS_ENABLED_TABLE_PROP)); + } + } } diff --git a/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/SnapshotsControllerTest.java b/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/SnapshotsControllerTest.java index 763daf0a2..7719eb799 100644 --- a/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/SnapshotsControllerTest.java +++ b/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/SnapshotsControllerTest.java @@ -20,6 +20,7 @@ import com.linkedin.openhouse.tables.model.TableDto; import com.linkedin.openhouse.tables.model.TableDtoPrimaryKey; import com.linkedin.openhouse.tables.repository.OpenHouseInternalRepository; +import com.linkedin.openhouse.tables.utils.AuthorizationUtils; import com.linkedin.openhouse.tablestest.annotation.CustomParameterResolver; import java.util.ArrayList; import java.util.Collections; @@ -312,6 +313,12 @@ public void testPutSnapshotsReplicaTableType(GetTableResponseBody getTableRespon @MethodSource("responseBodyFeeder") public void testPutSnapshotsReplaceCommit(GetTableResponseBody getTableResponseBody) throws Exception { + // RTAS is disabled by default; enable it via the table property so the replace commit is + // allowed. + Map propsWithRtas = new HashMap<>(getTableResponseBody.getTableProperties()); + propsWithRtas.put(AuthorizationUtils.RTAS_ENABLED_TABLE_PROP, "true"); + getTableResponseBody = getTableResponseBody.toBuilder().tableProperties(propsWithRtas).build(); + // Step 1: Create a table MvcResult createResult = RequestAndValidateHelper.createTableAndValidateResponse( diff --git a/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/TablesControllerTest.java b/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/TablesControllerTest.java index f44c1748b..0d9fc3f05 100644 --- a/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/TablesControllerTest.java +++ b/services/tables/src/test/java/com/linkedin/openhouse/tables/e2e/h2/TablesControllerTest.java @@ -47,6 +47,7 @@ import com.linkedin.openhouse.tables.repository.OpenHouseInternalRepository; import com.linkedin.openhouse.tables.toggle.model.TableToggleStatus; import com.linkedin.openhouse.tables.toggle.repository.ToggleStatusesRepository; +import com.linkedin.openhouse.tables.utils.AuthorizationUtils; import io.micrometer.core.instrument.search.MeterNotFoundException; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import java.util.ArrayList; @@ -829,8 +830,14 @@ public void testCreateTableWithIncorrectTableTypeThrowsException() throws Except @SneakyThrows @Test public void testStagedReplace() { - GetTableResponseBody stageReplaceTable = + GetTableResponseBody baseTable = TableModelConstants.buildGetTableResponseBodyWithDbTbl("d_sr", "t_sr"); + // RTAS is disabled by default; enable it via the table property so the staged replace is + // allowed. + Map propsWithRtas = new HashMap<>(baseTable.getTableProperties()); + propsWithRtas.put(AuthorizationUtils.RTAS_ENABLED_TABLE_PROP, "true"); + GetTableResponseBody stageReplaceTable = + baseTable.toBuilder().tableProperties(propsWithRtas).build(); // Create a real table first MvcResult createResult = @@ -904,6 +911,38 @@ public void testStagedReplace() { RequestAndValidateHelper.deleteTableAndValidateResponse(mvc, stageReplaceTable); } + @SneakyThrows + @Test + public void testStagedReplaceFailsWhenRtasDisabled() { + // Table created without enabling RTAS; replace is disabled by default. + GetTableResponseBody table = + TableModelConstants.buildGetTableResponseBodyWithDbTbl("d_sr_dis", "t_sr_dis"); + MvcResult createResult = + RequestAndValidateHelper.createTableAndValidateResponse(table, mvc, storageManager); + String originalTableLocation = + JsonPath.read(createResult.getResponse().getContentAsString(), "$.tableLocation"); + + // A stageReplace against a table without replaceEnabled should be rejected. + mvc.perform( + MockMvcRequestBuilders.post( + String.format( + ValidationUtilities.CURRENT_MAJOR_VERSION_PREFIX + "/databases/%s/tables/", + table.getDatabaseId())) + .contentType(MediaType.APPLICATION_JSON) + .content( + buildCreateUpdateTableRequestBody(table) + .toBuilder() + .baseTableVersion(originalTableLocation) + .stageReplace(true) + .build() + .toJson()) + .accept(MediaType.APPLICATION_JSON)) + .andExpect(status().isBadRequest()) + .andExpect(jsonPath("$.message", containsString("Replace (RTAS) is not enabled"))); + + RequestAndValidateHelper.deleteTableAndValidateResponse(mvc, table); + } + @SneakyThrows @Test public void testStagedCreateDoesntExistInConsecutiveCalls() { From db5bd713abe5a979ff2d5dbfef7c048f39dacb07 Mon Sep 17 00:00:00 2001 From: Levi Jiang Date: Wed, 24 Jun 2026 13:36:23 -0700 Subject: [PATCH 2/2] test(spark): add negative RTAS itest for disabled-by-default gate Add RTASTest.testRTASFailsWhenReplaceDisabled which creates a table without opting into RTAS and asserts that REPLACE TABLE ... AS SELECT is rejected with a BadRequestException carrying the "Replace (RTAS) is not enabled" message. Co-Authored-By: Claude Opus 4.8 --- .../openhouse/spark/catalogtest/RTASTest.java | 24 +++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java b/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java index ad921977e..1e593d8ef 100644 --- a/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java +++ b/integrations/spark/spark-3.1/openhouse-spark-itest/src/test/java/com/linkedin/openhouse/spark/catalogtest/RTASTest.java @@ -10,6 +10,7 @@ import org.apache.iceberg.Table; import org.apache.iceberg.catalog.Catalog; import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.iceberg.exceptions.BadRequestException; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.types.Types; import org.apache.spark.sql.Row; @@ -115,6 +116,29 @@ public void testRTAS() throws Exception { } } + @Test + public void testRTASFailsWhenReplaceDisabled() throws Exception { + try (SparkSession spark = getSparkSession()) { + // create the table without opting into RTAS; replace is disabled by default + spark.sql( + String.format( + "CREATE TABLE %s USING iceberg AS SELECT * FROM %s", tableName, sourceName)); + + // REPLACE TABLE should be rejected because 'replace.enabled' is not set on the table + BadRequestException exception = + assertThrows( + BadRequestException.class, + () -> + spark.sql( + String.format( + "REPLACE TABLE %s USING iceberg AS SELECT * FROM %s", + tableName, sourceName))); + assertTrue( + exception.getMessage().contains("Replace (RTAS) is not enabled"), + "Expected an RTAS-disabled error but got: " + exception.getMessage()); + } + } + @Test public void testCreateRTAS() throws Exception { try (SparkSession spark = getSparkSession()) {