From fca0f750943270f5133f3a6b22ab0a6e96d22b59 Mon Sep 17 00:00:00 2001 From: Ferenc Csaky Date: Fri, 31 Jul 2026 13:24:50 +0200 Subject: [PATCH] feat: Introduce shallow query engine --- .gitignore | 1 + .../java/com/datasqrl/config/EngineType.java | 1 + .../java/com/datasqrl/config/JdbcDialect.java | 6 +- .../QueryEngineConfigConverterImpl.java | 11 +- .../datasqrl/engine/database/QueryEngine.java | 3 +- .../AbstractJDBCShallowQueryEngine.java | 52 ++++ .../database/relational/DuckDBEngine.java | 7 +- .../relational/DuckDbStatementFactory.java | 12 +- .../relational/IcebergStatementFactory.java | 6 - .../relational/JdbcStatementFactory.java | 3 - .../relational/PostgresStatementFactory.java | 6 - .../database/relational/SnowflakeEngine.java | 15 +- .../relational/SnowflakeStatementFactory.java | 6 - .../engine/pipeline/SimplePipeline.java | 39 +-- .../java/com/datasqrl/util/EngineUtil.java | 22 +- .../com/datasqrl/EngineValidationTest.java | 33 ++- .../engine-validation/package-fail.json | 14 -- .../package-no-query-engine-fail.json | 7 + .../package-shallow-query-engine-fail.json | 7 + .../package-snowflake-no-server.json | 14 ++ .../engine-validation/package-snowflake.json | 14 ++ .../EngineValidationTest/package-fail.txt | 2 - .../package-no-query-engine-fail.txt | 2 + .../package-shallow-query-engine-fail.txt | 2 + .../package-snowflake-no-server.txt | 69 ++++++ .../package-snowflake.txt | 234 ++++++++++++++++++ 26 files changed, 479 insertions(+), 109 deletions(-) create mode 100644 sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCShallowQueryEngine.java delete mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-fail.json create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-no-query-engine-fail.json create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-shallow-query-engine-fail.json create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake-no-server.json create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake.json delete mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-fail.txt create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-no-query-engine-fail.txt create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-shallow-query-engine-fail.txt create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake-no-server.txt create mode 100644 sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake.txt diff --git a/.gitignore b/.gitignore index 1b63f85995..382495fc08 100644 --- a/.gitignore +++ b/.gitignore @@ -2,6 +2,7 @@ ## SQRL ############################## plan-output/ +sqrl_iceberg_data/ ############################## ## Java diff --git a/sqrl-planner/src/main/java/com/datasqrl/config/EngineType.java b/sqrl-planner/src/main/java/com/datasqrl/config/EngineType.java index d349e0a0b8..ab4e750ad9 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/config/EngineType.java +++ b/sqrl-planner/src/main/java/com/datasqrl/config/EngineType.java @@ -21,6 +21,7 @@ public enum EngineType { SERVER, LOG, QUERY, + SHALLOW_QUERY, EXPORT; public boolean isWrite() { diff --git a/sqrl-planner/src/main/java/com/datasqrl/config/JdbcDialect.java b/sqrl-planner/src/main/java/com/datasqrl/config/JdbcDialect.java index 98df9054f7..87ce18b0f9 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/config/JdbcDialect.java +++ b/sqrl-planner/src/main/java/com/datasqrl/config/JdbcDialect.java @@ -25,13 +25,11 @@ public enum JdbcDialect { SQLServer, H2, SQLite, - Iceberg, - Snowflake, - DuckDB; + Iceberg; private final String[] synonyms; - private JdbcDialect(String... synonyms) { + JdbcDialect(String... synonyms) { this.synonyms = synonyms; } diff --git a/sqrl-planner/src/main/java/com/datasqrl/config/QueryEngineConfigConverterImpl.java b/sqrl-planner/src/main/java/com/datasqrl/config/QueryEngineConfigConverterImpl.java index 1e44d32056..03b63748a6 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/config/QueryEngineConfigConverterImpl.java +++ b/sqrl-planner/src/main/java/com/datasqrl/config/QueryEngineConfigConverterImpl.java @@ -36,14 +36,19 @@ public List convertConfigsToJson() { for (var engine : enabledQueryEngines) { var queryEngine = (QueryEngine) engine; - var engineConf = packageJson.getEngines().getEngineConfig(queryEngine.getName()).get(); + var engineConf = packageJson.getEngines().getEngineConfig(queryEngine.getName()); + if (engineConf.isEmpty()) { + continue; + } - if (engineConf instanceof EngineConfigImpl impl) { + if (engineConf.get() instanceof EngineConfigImpl impl) { var engineConfigMap = impl.sqrlConfig.toMap(); var rootNode = JsonUtils.MAPPER.createObjectNode(); var configNode = JsonUtils.MAPPER.valueToTree(engineConfigMap); - rootNode.set(queryEngine.serverConfigName(), configNode); + queryEngine + .serverConfigName() + .ifPresent(engineConfigName -> rootNode.set(engineConfigName, configNode)); convertedConfigs.add(rootNode); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/QueryEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/QueryEngine.java index d72d1bcfe9..66702541c6 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/QueryEngine.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/QueryEngine.java @@ -18,6 +18,7 @@ import com.datasqrl.engine.EnginePhysicalPlan; import com.datasqrl.engine.ExecutionEngine; import com.datasqrl.planner.dag.plan.MaterializationStagePlan; +import java.util.Optional; /** * A {@link QueryEngine} executes queries against a {@link DatabaseEngine} that supports the query @@ -28,5 +29,5 @@ public interface QueryEngine extends ExecutionEngine { EnginePhysicalPlan plan(MaterializationStagePlan stagePlan); - String serverConfigName(); + Optional serverConfigName(); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCShallowQueryEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCShallowQueryEngine.java new file mode 100644 index 0000000000..f5e5640757 --- /dev/null +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCShallowQueryEngine.java @@ -0,0 +1,52 @@ +/* + * Copyright © 2021 DataSQRL (contact@datasqrl.com) + * + * Licensed 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 com.datasqrl.engine.database.relational; + +import static com.datasqrl.engine.EngineFeature.STANDARD_QUERY; + +import com.datasqrl.config.ConnectorFactoryFactory; +import com.datasqrl.config.EngineType; +import com.datasqrl.config.PackageJson.EngineConfig; +import com.datasqrl.engine.database.QueryEngine; +import com.datasqrl.planner.tables.FlinkTableBuilder; +import com.datasqrl.server.jdbc.DatabaseType; +import java.util.Optional; +import lombok.NonNull; + +/** Abstract implementation of a relational {@link QueryEngine}. */ +public abstract class AbstractJDBCShallowQueryEngine extends AbstractJDBCEngine + implements QueryEngine { + + protected AbstractJDBCShallowQueryEngine( + String name, @NonNull EngineConfig engineConfig, ConnectorFactoryFactory connectorFactory) { + super(name, EngineType.SHALLOW_QUERY, STANDARD_QUERY, engineConfig, connectorFactory); + } + + @Override + public final Optional serverConfigName() { + return Optional.empty(); + } + + @Override + protected final DatabaseType getDatabaseType() { + return DatabaseType.NONE; + } + + @Override + protected String getConnectorTableName(FlinkTableBuilder tableBuilder) { + throw new UnsupportedOperationException(); + } +} diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDBEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDBEngine.java index ae051b0fb3..7c6c6c70af 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDBEngine.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDBEngine.java @@ -20,6 +20,7 @@ import com.datasqrl.config.PackageJson; import com.datasqrl.server.jdbc.DatabaseType; import jakarta.inject.Inject; +import java.util.Optional; import lombok.NonNull; public class DuckDBEngine extends AbstractJDBCQueryEngine { @@ -33,8 +34,8 @@ public DuckDBEngine(@NonNull PackageJson json, ConnectorFactoryFactory connector } @Override - public String serverConfigName() { - return "duckDbConfig"; + public Optional serverConfigName() { + return Optional.of("duckDbConfig"); } @Override @@ -49,6 +50,6 @@ protected DatabaseType getDatabaseType() { @Override public JdbcStatementFactory getStatementFactory() { - return new DuckDbStatementFactory(engineConfig); + return new DuckDbStatementFactory(); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java index da581c76b2..0999860971 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java @@ -25,8 +25,6 @@ import com.datasqrl.calcite.Dialect; import com.datasqrl.calcite.dialect.DuckDbSqlDialect; import com.datasqrl.calcite.type.TypeFactory; -import com.datasqrl.config.JdbcDialect; -import com.datasqrl.config.PackageJson.EngineConfig; import com.datasqrl.engine.database.relational.ddl.GenericCreateTableDdlFactory; import com.datasqrl.plan.global.IndexDefinition; import com.datasqrl.planner.dag.plan.MaterializationStagePlan.Query; @@ -51,18 +49,10 @@ public class DuckDbStatementFactory extends AbstractJdbcStatementFactory { - private final EngineConfig engineConfig; - - public DuckDbStatementFactory(EngineConfig engineConfig) { + public DuckDbStatementFactory() { super( Dialect.DUCKDB, new GenericCreateTableDdlFactory()); // Iceberg creates the tables, DuckDB only queries - this.engineConfig = engineConfig; - } - - @Override - public JdbcDialect getDialect() { - return JdbcDialect.DuckDB; } @Override diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java index 955b02b94a..d5bc344bcc 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java @@ -18,7 +18,6 @@ import com.datasqrl.calcite.Dialect; import com.datasqrl.calcite.OperatorRuleTransformer; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; -import com.datasqrl.config.JdbcDialect; import com.datasqrl.engine.database.relational.ddl.IcebergCreateTableDdlFactory; import com.datasqrl.plan.global.IndexDefinition; import com.datasqrl.planner.hint.DataTypeHint; @@ -41,11 +40,6 @@ protected SqlDataTypeSpec getSqlType(RelDataType type, Optional hi return ExtendedPostgresSqlDialect.DEFAULT.getCastSpec(type); } - @Override - public JdbcDialect getDialect() { - return JdbcDialect.Postgres; - } - @Override public boolean supportsQueries() { return false; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatementFactory.java index fdd6ec7607..cc20ff78a9 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatementFactory.java @@ -15,7 +15,6 @@ */ package com.datasqrl.engine.database.relational; -import com.datasqrl.config.JdbcDialect; import com.datasqrl.plan.global.IndexDefinition; import com.datasqrl.planner.dag.plan.MaterializationStagePlan.Query; import java.util.Collection; @@ -24,8 +23,6 @@ public interface JdbcStatementFactory { - JdbcDialect getDialect(); - JdbcStatement createTable(JdbcEngineCreateTable createTable); default List applyTableExtensions(Collection tables) { diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java index 8e2b3e5302..598c82dad1 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java @@ -19,7 +19,6 @@ import com.datasqrl.calcite.Dialect; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; -import com.datasqrl.config.JdbcDialect; import com.datasqrl.config.PackageJson.EngineConfig; import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.deployment.model.JdbcStatementModel.Type; @@ -64,11 +63,6 @@ public PostgresStatementFactory(int partitionTtlDivisor) { this.partitionTtlDivisor = partitionTtlDivisor; } - @Override - public JdbcDialect getDialect() { - return JdbcDialect.Postgres; - } - @Override protected SqlDataTypeSpec getSqlType(RelDataType type, Optional hint) { SqlDataTypeSpec spec = ExtendedPostgresSqlDialect.DEFAULT.getCastSpec(type); diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeEngine.java index 86ae4f32e7..7318a18d45 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeEngine.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeEngine.java @@ -18,11 +18,10 @@ import com.datasqrl.config.ConnectorFactoryFactory; import com.datasqrl.config.JdbcDialect; import com.datasqrl.config.PackageJson; -import com.datasqrl.server.jdbc.DatabaseType; import jakarta.inject.Inject; import lombok.NonNull; -public class SnowflakeEngine extends AbstractJDBCQueryEngine { +public class SnowflakeEngine extends AbstractJDBCShallowQueryEngine { @Inject public SnowflakeEngine(@NonNull PackageJson json, ConnectorFactoryFactory connectorFactory) { @@ -32,19 +31,9 @@ public SnowflakeEngine(@NonNull PackageJson json, ConnectorFactoryFactory connec connectorFactory); } - @Override - public String serverConfigName() { - return "snowflakeConfig"; - } - @Override protected JdbcDialect getDialect() { - return JdbcDialect.Snowflake; - } - - @Override - protected DatabaseType getDatabaseType() { - return DatabaseType.SNOWFLAKE; + return JdbcDialect.Iceberg; } @Override diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java index e27f3d7656..a02a23e690 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java @@ -18,7 +18,6 @@ import com.datasqrl.calcite.Dialect; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; import com.datasqrl.calcite.dialect.snowflake.SqlCreateIcebergTableFromObjectStorage; -import com.datasqrl.config.JdbcDialect; import com.datasqrl.config.PackageJson.EngineConfig; import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.database.relational.ddl.GenericCreateTableDdlFactory; @@ -40,11 +39,6 @@ public SnowflakeStatementFactory(EngineConfig engineConfig) { this.engineConfig = engineConfig; } - @Override - public JdbcDialect getDialect() { - return JdbcDialect.Snowflake; - } - @Override public JdbcStatement createTable(JdbcEngineCreateTable createTable) { var tableName = createTable.tableName(); diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/pipeline/SimplePipeline.java b/sqrl-planner/src/main/java/com/datasqrl/engine/pipeline/SimplePipeline.java index f794c01808..e89a3cb1d6 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/pipeline/SimplePipeline.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/pipeline/SimplePipeline.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.stream.Collectors; /** * A simple pipeline that has a single stream, log, and server engine with support for multiple @@ -38,8 +39,8 @@ public record SimplePipeline( HashMultimap downstream) implements ExecutionPipeline { - private static final List AVAILABLE_QUERY_ENGINES = - EngineUtil.getAvailableQueryEngineNames(); + private static final String QUERY_ENGINE_NAMES = + EngineUtil.formatEngineNames(EngineUtil.getAvailableQueryEngines()); public static SimplePipeline of(Map engines, ErrorCollector errors) { var upstream = HashMultimap.create(); @@ -76,8 +77,28 @@ public static SimplePipeline of(Map engines, ErrorColle streamStage.ifPresent(ss -> upstream.put(dbStage, ss)); serverStage.ifPresent(vs -> downstream.put(dbStage, vs)); + // Make sure if server is present, then a non-view query engine is also present if (serverStage.isPresent() && dbStage.engine() instanceof AbstractJDBCTableFormatEngine) { - validatePipelineForQueryEngine(dbStage.name(), engines, errors); + var queryStages = getStage(EngineType.QUERY, engines); + var shallowQueryStages = getStage(EngineType.SHALLOW_QUERY, engines); + + if (queryStages.isEmpty() && !shallowQueryStages.isEmpty()) { + var shallowQueryEngines = + shallowQueryStages.stream() + .map(EngineStage::name) + .map(s -> '\'' + s + '\'') + .collect(Collectors.joining(", ")); + + errors.fatal( + "When '%s' is enabled as a server, '%s' cannot use shallow query engines (%s) to process server queries because they are not integrated at the database level. Available query engines: %s", + serverStage.get().name(), dbStage.name(), shallowQueryEngines, QUERY_ENGINE_NAMES); + } + + if (queryStages.isEmpty()) { + errors.fatal( + "When '%s' is enabled as a server, '%s' requires a query engine to process server queries, but none are listed under 'enabled-engines'. Available query engines: %s", + serverStage.get().name(), dbStage.name(), QUERY_ENGINE_NAMES); + } } } @@ -129,18 +150,6 @@ private static Optional getSingleStage( "Expected a single %s engine but found multiple: %s".formatted(engineType, engineList)); } - private static void validatePipelineForQueryEngine( - String tableFormatEngineName, Map engines, ErrorCollector errors) { - var queryStages = getStage(EngineType.QUERY, engines); - if (!queryStages.isEmpty()) { - return; - } - - errors.fatal( - "Engine '%s' requires a query engine, but none are listed under 'enabled-engines'. Available options: %s", - tableFormatEngineName, AVAILABLE_QUERY_ENGINES); - } - @Override public Set getUpStreamFrom(ExecutionStage stage) { Preconditions.checkArgument(upstream.containsKey(stage), "Invalid stage: %s", stage); diff --git a/sqrl-planner/src/main/java/com/datasqrl/util/EngineUtil.java b/sqrl-planner/src/main/java/com/datasqrl/util/EngineUtil.java index b668599757..87a3e0cb36 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/util/EngineUtil.java +++ b/sqrl-planner/src/main/java/com/datasqrl/util/EngineUtil.java @@ -17,10 +17,8 @@ import com.datasqrl.config.EngineFactory; import com.datasqrl.engine.IExecutionEngine; -import com.datasqrl.engine.database.DatabaseEngine; -import com.datasqrl.engine.database.QueryEngine; -import java.util.List; -import java.util.Set; +import com.datasqrl.engine.database.relational.AbstractJDBCQueryEngine; +import java.util.stream.Collectors; import java.util.stream.Stream; import lombok.AccessLevel; import lombok.NoArgsConstructor; @@ -28,17 +26,15 @@ @NoArgsConstructor(access = AccessLevel.PRIVATE) public final class EngineUtil { - public static List getAvailableDatabaseEngineNames(String... exclusions) { - var exclusionSet = Set.of(exclusions); - - return getChildEngineFactories(DatabaseEngine.class) - .map(EngineFactory::getEngineName) - .filter(name -> !exclusionSet.contains(name)) - .toList(); + public static Stream getAvailableQueryEngines() { + return getChildEngineFactories(AbstractJDBCQueryEngine.class); } - public static List getAvailableQueryEngineNames() { - return getChildEngineFactories(QueryEngine.class).map(EngineFactory::getEngineName).toList(); + public static String formatEngineNames(Stream engines) { + return engines + .map(EngineFactory::getEngineName) + .map(name -> '\'' + name + '\'') + .collect(Collectors.joining(", ")); } public static Stream getChildEngineFactories( diff --git a/sqrl-testing/sqrl-testing-integration/src/test/java/com/datasqrl/EngineValidationTest.java b/sqrl-testing/sqrl-testing-integration/src/test/java/com/datasqrl/EngineValidationTest.java index 8cdb127ed6..a89f20deee 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/java/com/datasqrl/EngineValidationTest.java +++ b/sqrl-testing/sqrl-testing-integration/src/test/java/com/datasqrl/EngineValidationTest.java @@ -20,10 +20,13 @@ import static org.assertj.core.api.Assertions.assertThat; import com.datasqrl.SnapshotTestSupport.TestNameModifier; +import com.datasqrl.util.ArgumentsProviders; import com.datasqrl.util.SnapshotTest.Snapshot; import java.nio.file.Path; -import org.junit.jupiter.api.Test; +import java.util.function.Predicate; import org.junit.jupiter.api.extension.RegisterExtension; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ArgumentsSource; class EngineValidationTest { @@ -33,28 +36,40 @@ class EngineValidationTest { final CliCompileTestExtension snapshotExtension = new CliCompileTestExtension(Path.of("plan-output")); - @Test - void testInvalidEngine() { - var pkg = PROJECT_DIR.resolve("package-fail.json"); - assertThat(pkg).isRegularFile(); + @ParameterizedTest + @ArgumentsSource(ProjectCaseFiles.class) + void testEngineValidation(Path packageFile) { + assertThat(packageFile).isRegularFile(); - var testModifier = TestNameModifier.of(pkg); + var testModifier = TestNameModifier.of(packageFile); var expectFailure = testModifier == TestNameModifier.fail; var printMessages = testModifier == TestNameModifier.fail || testModifier == TestNameModifier.warn; - snapshotExtension.setSnapshot(Snapshot.of(getDisplayName(pkg), getClass())); + snapshotExtension.setSnapshot(Snapshot.of(getDisplayName(packageFile), getClass())); var hook = snapshotExtension.execute( PROJECT_DIR, "compile", - pkg.getFileName().toString(), + packageFile.getFileName().toString(), "-t", snapshotExtension.getOutputDir().toAbsolutePath().toString()); assertThat(hook.isFailed()).as(hook.getMessages()).isEqualTo(expectFailure); if (printMessages) { snapshotExtension.createMessageSnapshot(hook.getMessages()); } else { - snapshotExtension.createSnapshot(buildDir -> false, outputDir -> false, planDir -> true); + snapshotExtension.createSnapshot(buildDir -> false, getOutputDirFilter(), planDir -> true); + } + } + + private Predicate getOutputDirFilter() { + return path -> + path.getFileName().toString().contains("iceberg") + || path.getFileName().toString().equals("vertx.json"); + } + + static class ProjectCaseFiles extends ArgumentsProviders.PackageProvider { + ProjectCaseFiles() { + super(PROJECT_DIR); } } } diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-fail.json b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-fail.json deleted file mode 100644 index b15fa24d9f..0000000000 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-fail.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "version": "1", - "enabled-engines": ["flink", "iceberg", "vertx"], - "script": { - "main": "script.sqrl" - }, - "connectors" : { - "iceberg" : { - "warehouse":"warehouse", - "catalog-type":"hadoop", - "catalog-name": "mydatabase" - } - } -} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-no-query-engine-fail.json b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-no-query-engine-fail.json new file mode 100644 index 0000000000..0c5666a840 --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-no-query-engine-fail.json @@ -0,0 +1,7 @@ +{ + "version": "1", + "enabled-engines": ["flink", "iceberg", "vertx"], + "script": { + "main": "script.sqrl" + } +} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-shallow-query-engine-fail.json b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-shallow-query-engine-fail.json new file mode 100644 index 0000000000..129aee00d4 --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-shallow-query-engine-fail.json @@ -0,0 +1,7 @@ +{ + "version": "1", + "enabled-engines": ["flink", "iceberg", "vertx", "snowflake"], + "script": { + "main": "script.sqrl" + } +} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake-no-server.json b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake-no-server.json new file mode 100644 index 0000000000..3aff305964 --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake-no-server.json @@ -0,0 +1,14 @@ +{ + "version": "1", + "enabled-engines": ["flink", "iceberg", "snowflake"], + "script": { + "main": "script.sqrl" + }, + "engines": { + "snowflake": { + "catalog-name": "MyCatalog", + "external-volume": "MyNewVolume", + "url": "${SNOWFLAKE_JDBC_URL}" + } + } +} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake.json b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake.json new file mode 100644 index 0000000000..0e46ea95ef --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/engine-validation/package-snowflake.json @@ -0,0 +1,14 @@ +{ + "version": "1", + "enabled-engines": ["flink", "iceberg", "vertx", "duckdb", "snowflake"], + "script": { + "main": "script.sqrl" + }, + "engines": { + "snowflake": { + "catalog-name": "MyCatalog", + "external-volume": "MyNewVolume", + "url": "${SNOWFLAKE_JDBC_URL}" + } + } +} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-fail.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-fail.txt deleted file mode 100644 index 5ee2ca79fc..0000000000 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-fail.txt +++ /dev/null @@ -1,2 +0,0 @@ -[FATAL] Engine 'iceberg' requires a query engine, but none are listed under 'enabled-engines'. Available options: [duckdb, snowflake] - diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-no-query-engine-fail.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-no-query-engine-fail.txt new file mode 100644 index 0000000000..349626c905 --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-no-query-engine-fail.txt @@ -0,0 +1,2 @@ +[FATAL] When 'vertx' is enabled as a server, 'iceberg' requires a query engine to process server queries, but none are listed under 'enabled-engines'. Available query engines: 'duckdb' + diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-shallow-query-engine-fail.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-shallow-query-engine-fail.txt new file mode 100644 index 0000000000..13b477f2e3 --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-shallow-query-engine-fail.txt @@ -0,0 +1,2 @@ +[FATAL] When 'vertx' is enabled as a server, 'iceberg' cannot use shallow query engines ('snowflake') to process server queries because they are not integrated at the database level. Available query engines: 'duckdb' + diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake-no-server.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake-no-server.txt new file mode 100644 index 0000000000..072f2a6e0e --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake-no-server.txt @@ -0,0 +1,69 @@ +>>>iceberg-schema.sql +CREATE TABLE IF NOT EXISTS "Numbers" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")); +CREATE TABLE IF NOT EXISTS "Numbers2" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")) +>>>iceberg-snowflake-schema.sql +CREATE OR REPLACE ICEBERG TABLE Numbers EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers'; +CREATE OR REPLACE ICEBERG TABLE Numbers2 EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers2' +>>>iceberg.json +{ + "plans" : { + "" : { + "statements" : [ + { + "name" : "Numbers", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + }, + { + "name" : "Numbers2", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers2\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + } + ], + "standaloneExtensionStatements" : [ ] + }, + "snowflake" : { + "statements" : [ + { + "name" : "Numbers", + "type" : "TABLE", + "sql" : "CREATE OR REPLACE ICEBERG TABLE Numbers EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers'" + }, + { + "name" : "Numbers2", + "type" : "TABLE", + "sql" : "CREATE OR REPLACE ICEBERG TABLE Numbers2 EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers2'" + } + ], + "standaloneExtensionStatements" : [ ] + } + } +} diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake.txt new file mode 100644 index 0000000000..bd57b79cee --- /dev/null +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/EngineValidationTest/package-snowflake.txt @@ -0,0 +1,234 @@ +>>>iceberg-duckdb-schema.sql +CREATE TABLE IF NOT EXISTS "Numbers" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")); +CREATE TABLE IF NOT EXISTS "Numbers2" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")) +>>>iceberg-schema.sql +CREATE TABLE IF NOT EXISTS "Numbers" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")); +CREATE TABLE IF NOT EXISTS "Numbers2" ("id" INTEGER NOT NULL, PRIMARY KEY ("id")) +>>>iceberg-snowflake-schema.sql +CREATE OR REPLACE ICEBERG TABLE Numbers EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers'; +CREATE OR REPLACE ICEBERG TABLE Numbers2 EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers2' +>>>iceberg.json +{ + "plans" : { + "" : { + "statements" : [ + { + "name" : "Numbers", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + }, + { + "name" : "Numbers2", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers2\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + } + ], + "standaloneExtensionStatements" : [ ] + }, + "duckdb" : { + "statements" : [ + { + "name" : "Numbers", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + }, + { + "name" : "Numbers2", + "type" : "TABLE", + "sql" : "CREATE TABLE IF NOT EXISTS \"Numbers2\" (\"id\" INTEGER NOT NULL, PRIMARY KEY (\"id\"))", + "fields" : [ + { + "name" : "id", + "type" : "INTEGER", + "nullable" : false + } + ], + "primaryKey" : [ + "id" + ], + "partitionKey" : [ ], + "partitionType" : "NONE", + "numPartitions" : 0, + "ttl" : 0.0 + } + ], + "standaloneExtensionStatements" : [ ] + }, + "snowflake" : { + "statements" : [ + { + "name" : "Numbers", + "type" : "TABLE", + "sql" : "CREATE OR REPLACE ICEBERG TABLE Numbers EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers'" + }, + { + "name" : "Numbers2", + "type" : "TABLE", + "sql" : "CREATE OR REPLACE ICEBERG TABLE Numbers2 EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = 'Numbers2'" + } + ], + "standaloneExtensionStatements" : [ ] + } + } +} +>>>vertx.json +{ + "models" : { + "v1" : { + "queries" : [ + { + "type" : "args", + "parentType" : "Query", + "fieldName" : "Numbers", + "exec" : { + "arguments" : [ + { + "type" : "variable", + "path" : "limit" + }, + { + "type" : "variable", + "path" : "offset" + } + ], + "query" : { + "type" : "SqlQuery", + "sql" : "SELECT *\nFROM \"iceberg_scan\"('sqrl_iceberg_data/default_database/Numbers', ALLOW_MOVED_PATHS = TRUE)", + "parameters" : [ ], + "pagination" : "LIMIT_AND_OFFSET", + "cacheDurationMs" : 0, + "database" : "DUCKDB" + } + } + }, + { + "type" : "args", + "parentType" : "Query", + "fieldName" : "Numbers2", + "exec" : { + "arguments" : [ + { + "type" : "variable", + "path" : "limit" + }, + { + "type" : "variable", + "path" : "offset" + } + ], + "query" : { + "type" : "SqlQuery", + "sql" : "SELECT *\nFROM \"iceberg_scan\"('sqrl_iceberg_data/default_database/Numbers2', ALLOW_MOVED_PATHS = TRUE)", + "parameters" : [ ], + "pagination" : "LIMIT_AND_OFFSET", + "cacheDurationMs" : 0, + "database" : "DUCKDB" + } + } + } + ], + "mutations" : [ ], + "subscriptions" : [ ], + "operations" : [ + { + "function" : { + "name" : "GetNumbers", + "parameters" : { + "type" : "object", + "properties" : { + "offset" : { + "type" : "integer" + }, + "limit" : { + "type" : "integer" + } + }, + "required" : [ ] + } + }, + "format" : "JSON", + "apiQuery" : { + "query" : "query Numbers($limit: Int = 10, $offset: Int = 0) {\nNumbers(limit: $limit, offset: $offset) {\nid\n}\n\n}", + "queryName" : "Numbers", + "operationType" : "QUERY" + }, + "mcpMethod" : "TOOL", + "restMethod" : "GET", + "uriTemplate" : "queries/Numbers{?offset,limit}" + }, + { + "function" : { + "name" : "GetNumbers2", + "parameters" : { + "type" : "object", + "properties" : { + "offset" : { + "type" : "integer" + }, + "limit" : { + "type" : "integer" + } + }, + "required" : [ ] + } + }, + "format" : "JSON", + "apiQuery" : { + "query" : "query Numbers2($limit: Int = 10, $offset: Int = 0) {\nNumbers2(limit: $limit, offset: $offset) {\nid\n}\n\n}", + "queryName" : "Numbers2", + "operationType" : "QUERY" + }, + "mcpMethod" : "TOOL", + "restMethod" : "GET", + "uriTemplate" : "queries/Numbers2{?offset,limit}" + } + ], + "schema" : { + "type" : "string", + "schema" : "\"An RFC-3339 compliant Full Date Scalar\"\nscalar Date\n\n\"A DateTime scalar that handles both full RFC3339 and shorter timestamp formats\"\nscalar DateTime\n\n\"A JSON scalar\"\nscalar JSON\n\n\"24-hour clock time value string in the format `hh:mm:ss` or `hh:mm:ss.sss`.\"\nscalar LocalTime\n\n\"A 64-bit signed integer\"\nscalar Long\n\ntype Numbers {\n id: Int!\n}\n\ntype Numbers2 {\n id: Int!\n}\n\ntype Query {\n Numbers(limit: Int = 10, offset: Int = 0): [Numbers!]\n Numbers2(limit: Int = 10, offset: Int = 0): [Numbers2!]\n}\n\nenum _McpMethodType {\n NONE\n TOOL\n RESOURCE\n}\n\nenum _RestMethodType {\n NONE\n GET\n POST\n}\n\ndirective @api(mcp: _McpMethodType, rest: _RestMethodType, uri: String) on QUERY | MUTATION | FIELD_DEFINITION\n" + } + } + } +}