diff --git a/pom.xml b/pom.xml index 5b0ba56db2..86b33652cd 100644 --- a/pom.xml +++ b/pom.xml @@ -69,6 +69,7 @@ sqrl-cli sqrl-discovery sqrl-planner + sqrl-deployment-model sqrl-server sqrl-testing @@ -220,6 +221,11 @@ + + com.datasqrl + sqrl-deployment-model + ${project.version} + com.datasqrl sqrl-server-core diff --git a/sqrl-cli/src/main/java/com/datasqrl/cli/DatasqrlRun.java b/sqrl-cli/src/main/java/com/datasqrl/cli/DatasqrlRun.java index c4e4547be8..2179a4dd9b 100644 --- a/sqrl-cli/src/main/java/com/datasqrl/cli/DatasqrlRun.java +++ b/sqrl-cli/src/main/java/com/datasqrl/cli/DatasqrlRun.java @@ -21,6 +21,7 @@ import static com.datasqrl.env.EnvVariableNames.POSTGRES_USERNAME; import com.datasqrl.config.PackageJson; +import com.datasqrl.deployment.model.KafkaNewTopicModel; import com.datasqrl.engine.server.VertxEngineFactory; import com.datasqrl.flinkrunner.SqrlRunner; import com.datasqrl.flinkrunner.utils.EnvUtils; @@ -260,7 +261,7 @@ private void initKafka() { var topicsToCreate = new HashSet(); Stream.concat(kafkaPlan.topics().stream(), kafkaPlan.testRunnerTopics().stream()) - .map(com.datasqrl.engine.log.kafka.NewTopic::topicName) + .map(KafkaNewTopicModel::topicName) .forEach(topicsToCreate::add); var bootstrapServers = getenv(KAFKA_BOOTSTRAP_SERVERS); @@ -304,9 +305,9 @@ private void initPostgres() { DriverManager.getConnection( getenv(POSTGRES_JDBC_URL), getenv(POSTGRES_USERNAME), getenv(POSTGRES_PASSWORD))) { for (var jdbcStmt : statements) { - log.info("Executing statement {} of type {}", jdbcStmt.getName(), jdbcStmt.getType()); + log.info("Executing statement {} of type {}", jdbcStmt.name(), jdbcStmt.type()); try (Statement stmt = connection.createStatement()) { - stmt.execute(jdbcStmt.getSql()); + stmt.execute(jdbcStmt.sql()); } catch (Exception e) { e.printStackTrace(); assert false : e.getMessage(); diff --git a/sqrl-cli/src/main/java/com/datasqrl/compile/CompilationProcess.java b/sqrl-cli/src/main/java/com/datasqrl/compile/CompilationProcess.java index 5d6134015f..9b190e0964 100644 --- a/sqrl-cli/src/main/java/com/datasqrl/compile/CompilationProcess.java +++ b/sqrl-cli/src/main/java/com/datasqrl/compile/CompilationProcess.java @@ -18,9 +18,9 @@ import com.datasqrl.config.GraphqlSourceLoader; import com.datasqrl.config.PackageJson; import com.datasqrl.config.WorkspacePaths; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.PhysicalPlan; import com.datasqrl.engine.database.relational.JdbcPhysicalPlan; -import com.datasqrl.engine.database.relational.JdbcStatement; import com.datasqrl.engine.server.ServerPhysicalPlan; import com.datasqrl.engine.stream.flink.FlinkStreamEngine; import com.datasqrl.error.ErrorCode; @@ -31,6 +31,7 @@ import com.datasqrl.planner.SqlScriptPlanner; import com.datasqrl.planner.Sqrl2FlinkSQLTranslator; import com.datasqrl.planner.dag.DAGPlanner; +import com.datasqrl.planner.dag.plan.MutationDatabase; import com.datasqrl.server.GenerateServerModel; import com.datasqrl.util.ServiceLoaderDiscovery; import java.nio.file.Path; @@ -111,7 +112,7 @@ public Pair executeCompilation(Optional testsPath) var jdbcViews = physicalPlan .getPlans(JdbcPhysicalPlan.class) - .map(p -> p.getStatementsForType(JdbcStatement.Type.VIEW)) + .map(p -> p.getStatementsForType(Type.VIEW)) .findFirst() .orElse(List.of()); @@ -126,7 +127,8 @@ public Pair executeCompilation(Optional testsPath) .ifPresent( compareDb -> errors.checkFatal( - mutationDatabase.isBackwardsCompatible(compareDb, environment, errors), + MutationDatabase.isBackwardsCompatible( + mutationDatabase, compareDb, environment, errors), "The mutation tables defined by the script are not backwards compatible with the provided database. See warnings above for incompatibilities.")); return Pair.of(physicalPlan, testPlan); } diff --git a/sqrl-cli/src/main/java/com/datasqrl/compile/DagWriter.java b/sqrl-cli/src/main/java/com/datasqrl/compile/DagWriter.java index e5ce7f1b2d..77ce17080e 100644 --- a/sqrl-cli/src/main/java/com/datasqrl/compile/DagWriter.java +++ b/sqrl-cli/src/main/java/com/datasqrl/compile/DagWriter.java @@ -18,9 +18,9 @@ import com.datasqrl.config.PackageJson.CompilerConfig; import com.datasqrl.config.PackageJson.ExplainConfig; import com.datasqrl.config.WorkspacePaths; +import com.datasqrl.deployment.model.MutationDatabaseModel; import com.datasqrl.plan.global.PipelineDAGExporter; import com.datasqrl.planner.dag.PipelineDAG; -import com.datasqrl.planner.dag.plan.MutationDatabase; import com.datasqrl.serializer.Deserializer; import com.google.common.io.Resources; import java.io.IOException; @@ -55,12 +55,12 @@ public class DagWriter { private final CompilerConfig compilerConfig; @SneakyThrows - void run(PipelineDAG dag, String source, MutationDatabase mutationDatabase) { + void run(PipelineDAG dag, String source, MutationDatabaseModel mutationDatabaseModel) { writeExplain(dag); writeFile(workspacePaths.buildDir().resolve(FULL_SOURCE_FILENAME), source); writeFile( workspacePaths.buildDir().resolve(DATABASE_FILENAME), - Deserializer.INSTANCE.toJson(mutationDatabase)); + Deserializer.INSTANCE.toJson(mutationDatabaseModel)); } void writeInferredSchema(String inferredSchema) { diff --git a/sqrl-cli/src/test/java/com/datasqrl/cli/DatasqrlTestTest.java b/sqrl-cli/src/test/java/com/datasqrl/cli/DatasqrlTestTest.java index a22df41043..3ecef2d65a 100644 --- a/sqrl-cli/src/test/java/com/datasqrl/cli/DatasqrlTestTest.java +++ b/sqrl-cli/src/test/java/com/datasqrl/cli/DatasqrlTestTest.java @@ -23,7 +23,11 @@ import com.datasqrl.cli.output.NoOutputFormatter; import com.datasqrl.cli.output.TestOutputManager; +import com.datasqrl.compile.TestPlan; import com.datasqrl.config.PackageJson; +import com.datasqrl.deployment.model.JdbcStatementModel; +import com.datasqrl.engine.database.relational.GenericJdbcStatement; +import com.datasqrl.util.SqrlObjectMapper; import java.nio.file.Files; import java.nio.file.Path; import java.util.Map; @@ -68,4 +72,28 @@ void run_whenPipelineFailsToStart_recordsFailureWithNonZeroExit() throws Excepti var exitCode = underTest.run(); assertThat(exitCode).isNotZero(); } + + @Test + void givenTestPlanWithJdbcViews_whenDeserialized_thenUsesGenericJdbcStatements() + throws Exception { + var json = + """ + { + "jdbcViews": [{ + "name": "orders_view", + "type": "VIEW", + "sql": "CREATE VIEW orders_view AS SELECT 1" + }], + "queries": [], + "mutations": [], + "subscriptions": [] + } + """; + + var plan = SqrlObjectMapper.INSTANCE.readValue(json, TestPlan.class); + + assertThat(plan.getJdbcViews()).hasSize(1); + assertThat(plan.getJdbcViews().get(0)).isInstanceOf(GenericJdbcStatement.class); + assertThat(plan.getJdbcViews().get(0).getType()).isEqualTo(JdbcStatementModel.Type.VIEW); + } } diff --git a/sqrl-cli/src/test/java/com/datasqrl/util/OsProcessManagerTest.java b/sqrl-cli/src/test/java/com/datasqrl/util/OsProcessManagerTest.java index 5ed5bb7786..13070cfd24 100644 --- a/sqrl-cli/src/test/java/com/datasqrl/util/OsProcessManagerTest.java +++ b/sqrl-cli/src/test/java/com/datasqrl/util/OsProcessManagerTest.java @@ -22,11 +22,10 @@ import static org.mockito.Mockito.mockStatic; import static org.mockito.Mockito.when; -import com.datasqrl.engine.database.relational.GenericJdbcStatement; -import com.datasqrl.engine.database.relational.JdbcPhysicalPlan; -import com.datasqrl.engine.database.relational.JdbcStatement; -import com.datasqrl.engine.log.kafka.KafkaPhysicalPlan; -import com.datasqrl.engine.log.kafka.NewTopic; +import com.datasqrl.deployment.model.JdbcPlanModel; +import com.datasqrl.deployment.model.JdbcStatementModel; +import com.datasqrl.deployment.model.KafkaNewTopicModel; +import com.datasqrl.deployment.model.KafkaPlanModel; import com.datasqrl.env.GlobalEnvironmentStore; import java.io.File; import java.io.IOException; @@ -78,7 +77,7 @@ void givenLogFileExists_whenReadServiceLogFile_thenReturnsContent() throws Excep filesMocked.when(() -> Files.exists(mockLogFile)).thenReturn(true); filesMocked .when(() -> Files.readAllLines(mockLogFile)) - .thenReturn(java.util.List.of("Line 1", "Line 2", "Line 3")); + .thenReturn(List.of("Line 1", "Line 2", "Line 3")); // When String result = serviceManager.readServiceLogFile(serviceName); @@ -147,7 +146,7 @@ void givenCustomEnvironmentVariables_whenStartDependentServices_thenSetsSystemPr // Mock that no services are needed configMocked .when(() -> ConfigLoaderUtils.loadKafkaPhysicalPlan(mockPlanDir)) - .thenReturn(Optional.of(new KafkaPhysicalPlan(List.of(), List.of()))); + .thenReturn(Optional.of(new KafkaPlanModel(List.of(), List.of()))); configMocked .when(() -> ConfigLoaderUtils.loadPostgresPhysicalPlan(mockPlanDir)) .thenReturn(Optional.empty()); @@ -342,10 +341,10 @@ void givenPlanDirWithKafkaTopics_whenStartDependentServices_thenStartsRedpanda() filesMocked.when(() -> Files.createDirectories(any(Path.class))).thenReturn(mockPath); // Mock that Kafka topics are found but no Postgres statements - var mockTopic = mock(NewTopic.class); + var mockTopic = mock(KafkaNewTopicModel.class); configMocked .when(() -> ConfigLoaderUtils.loadKafkaPhysicalPlan(mockPlanDir)) - .thenReturn(Optional.of(new KafkaPhysicalPlan(List.of(mockTopic), List.of()))); + .thenReturn(Optional.of(new KafkaPlanModel(List.of(mockTopic), List.of()))); configMocked .when(() -> ConfigLoaderUtils.loadPostgresPhysicalPlan(mockPlanDir)) .thenReturn(Optional.empty()); @@ -396,8 +395,9 @@ void givenPlanDirWithPostgresStatements_whenStartDependentServices_thenStartsPos .when(() -> ConfigLoaderUtils.loadKafkaPhysicalPlan(mockPlanDir)) .thenReturn(Optional.empty()); var mockStatement = - new GenericJdbcStatement("test", JdbcStatement.Type.TABLE, "CREATE TABLE test"); - var mockJdbcPlan = JdbcPhysicalPlan.builder().statement(mockStatement).build(); + new JdbcStatementModel( + "test", JdbcStatementModel.Type.TABLE, "CREATE TABLE test", null, null); + var mockJdbcPlan = new JdbcPlanModel(List.of(mockStatement), List.of()); configMocked .when(() -> ConfigLoaderUtils.loadPostgresPhysicalPlan(mockPlanDir)) .thenReturn(Optional.of(mockJdbcPlan)); @@ -442,13 +442,14 @@ void givenPlanDirWithBothServices_whenStartDependentServices_thenStartsBothServi filesMocked.when(() -> Files.list(any(Path.class))).thenReturn(Stream.of(mockPath)); // Mock that both Kafka topics and Postgres statements are found - var mockTopic = mock(NewTopic.class); + var mockTopic = mock(KafkaNewTopicModel.class); configMocked .when(() -> ConfigLoaderUtils.loadKafkaPhysicalPlan(mockPlanDir)) - .thenReturn(Optional.of(new KafkaPhysicalPlan(List.of(mockTopic), List.of()))); + .thenReturn(Optional.of(new KafkaPlanModel(List.of(mockTopic), List.of()))); var mockStatement = - new GenericJdbcStatement("test", JdbcStatement.Type.TABLE, "CREATE TABLE test"); - var mockJdbcPlan = JdbcPhysicalPlan.builder().statement(mockStatement).build(); + new JdbcStatementModel( + "test", JdbcStatementModel.Type.TABLE, "CREATE TABLE test", null, null); + var mockJdbcPlan = new JdbcPlanModel(List.of(mockStatement), List.of()); configMocked .when(() -> ConfigLoaderUtils.loadPostgresPhysicalPlan(mockPlanDir)) .thenReturn(Optional.of(mockJdbcPlan)); diff --git a/sqrl-deployment-model/pom.xml b/sqrl-deployment-model/pom.xml new file mode 100644 index 0000000000..95b11d73cb --- /dev/null +++ b/sqrl-deployment-model/pom.xml @@ -0,0 +1,36 @@ + + + + 4.0.0 + + com.datasqrl + sqrl-root + 1.0-SNAPSHOT + + + sqrl-deployment-model + SQRL :: Deployment Model + + + + com.fasterxml.jackson.core + jackson-annotations + + + diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/stream/StreamPhysicalPlan.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/FlinkPlanModel.java similarity index 69% rename from sqrl-planner/src/main/java/com/datasqrl/engine/stream/StreamPhysicalPlan.java rename to sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/FlinkPlanModel.java index cc9c738289..c7fac06a69 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/stream/StreamPhysicalPlan.java +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/FlinkPlanModel.java @@ -13,8 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.datasqrl.engine.stream; +package com.datasqrl.deployment.model; -import com.datasqrl.engine.EnginePhysicalPlan; +import java.util.List; +import java.util.Set; -public interface StreamPhysicalPlan extends EnginePhysicalPlan {} +/** The contents of the {@code flink.json} deployment file. */ +public record FlinkPlanModel( + List flinkSql, Set connectors, Set formats, Set functions) {} diff --git a/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcPlanModel.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcPlanModel.java new file mode 100644 index 0000000000..2ec511973a --- /dev/null +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcPlanModel.java @@ -0,0 +1,33 @@ +/* + * 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.deployment.model; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import java.util.List; + +/** The contents of a JDBC database plan file such as {@code postgres.json}. */ +@JsonIgnoreProperties(ignoreUnknown = true) +public record JdbcPlanModel( + List statements, List standaloneExtensionStatements) { + + public JdbcPlanModel { + statements = statements == null ? List.of() : List.copyOf(statements); + standaloneExtensionStatements = + standaloneExtensionStatements == null + ? List.of() + : List.copyOf(standaloneExtensionStatements); + } +} diff --git a/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcStatementModel.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcStatementModel.java new file mode 100644 index 0000000000..1be6854311 --- /dev/null +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/JdbcStatementModel.java @@ -0,0 +1,59 @@ +/* + * 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.deployment.model; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonInclude; +import java.time.Duration; +import java.util.List; + +/** A rendered statement in a JDBC deployment file. */ +@JsonIgnoreProperties(ignoreUnknown = true) +@JsonInclude(JsonInclude.Include.NON_NULL) +public record JdbcStatementModel( + String name, + Type type, + String sql, + String description, + List fields, + List primaryKey, + List partitionKey, + PartitionType partitionType, + Integer numPartitions, + Duration ttl) { + + public JdbcStatementModel( + String name, Type type, String sql, String description, List fields) { + this(name, type, sql, description, fields, null, null, null, null, null); + } + + public enum Type { + TABLE, + VIEW, + QUERY, + INDEX, + EXTENSION + } + + public enum PartitionType { + NONE, + HASH, + LIST, + RANGE + } + + public record Field(String name, String type, boolean nullable, String description) {} +} diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/NewTopic.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaNewTopicModel.java similarity index 71% rename from sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/NewTopic.java rename to sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaNewTopicModel.java index adc47bebac..8677ddc59c 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/NewTopic.java +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaNewTopicModel.java @@ -13,13 +13,13 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.datasqrl.engine.log.kafka; +package com.datasqrl.deployment.model; -import com.datasqrl.engine.database.EngineCreateTable; import java.util.List; import java.util.Map; -public record NewTopic( +/** A Kafka topic definition in a deployment file. */ +public record KafkaNewTopicModel( String topicName, String tableName, String format, @@ -28,10 +28,10 @@ public record NewTopic( Type type, List messageKeys, String messageSchema, - Map config) - implements EngineCreateTable { + Map config) { - public NewTopic(String topicName, String tableName, int numPartitions, short replicationFactor) { + public KafkaNewTopicModel( + String topicName, String tableName, int numPartitions, short replicationFactor) { this( topicName, tableName, @@ -44,8 +44,8 @@ public NewTopic(String topicName, String tableName, int numPartitions, short rep Map.of()); } - public NewTopic(String topicName, String tableName) { - this(topicName, tableName, null, 1, (short) 1, Type.SUBSCRIPTION, List.of(), "", Map.of()); + public KafkaNewTopicModel(String topicName, String tableName) { + this(topicName, tableName, 1, (short) 1); } public enum Type { diff --git a/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaPlanModel.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaPlanModel.java new file mode 100644 index 0000000000..16e14dcbce --- /dev/null +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/KafkaPlanModel.java @@ -0,0 +1,36 @@ +/* + * 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.deployment.model; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import java.util.List; + +/** The contents of a Kafka deployment file. */ +@JsonIgnoreProperties(ignoreUnknown = true) +public record KafkaPlanModel( + List topics, List testRunnerTopics) { + + public KafkaPlanModel { + topics = topics == null ? List.of() : List.copyOf(topics); + testRunnerTopics = testRunnerTopics == null ? List.of() : List.copyOf(testRunnerTopics); + } + + @JsonIgnore + public boolean isEmpty() { + return topics.isEmpty() && testRunnerTopics.isEmpty(); + } +} diff --git a/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/MutationDatabaseModel.java b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/MutationDatabaseModel.java new file mode 100644 index 0000000000..607032ccde --- /dev/null +++ b/sqrl-deployment-model/src/main/java/com/datasqrl/deployment/model/MutationDatabaseModel.java @@ -0,0 +1,38 @@ +/* + * 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.deployment.model; + +import com.fasterxml.jackson.annotation.JsonInclude; +import java.util.List; +import java.util.Map; + +/** The contents of the {@code pipeline_mutation_database.json} deployment file. */ +@JsonInclude(JsonInclude.Include.NON_EMPTY) +public record MutationDatabaseModel(List tables) { + + public record Table( + String canonicalName, + String engine, + String createTableSql, + TableDefinition definition, + Map configOptions, + String documentation) {} + + public record TableDefinition( + List columns, List primaryKey, List partitionKey) {} + + public record ColumnDefinition(String name, String spec, String documentation) {} +} diff --git a/sqrl-planner/pom.xml b/sqrl-planner/pom.xml index d3b1673bef..3488d67f4a 100644 --- a/sqrl-planner/pom.xml +++ b/sqrl-planner/pom.xml @@ -28,6 +28,16 @@ SQRL :: Planner + + com.datasqrl + sqrl-deployment-model + + + + com.datasqrl + sqrl-server-vertx-base + + org.apache.avro avro @@ -114,11 +124,7 @@ com.h2database h2 - - com.datasqrl - sqrl-server-vertx-base - compile - + org.apache.flink flink-test-utils diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/EnginePhysicalPlan.java b/sqrl-planner/src/main/java/com/datasqrl/engine/EnginePhysicalPlan.java index 81b6eba6b5..ed17adb9a5 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/EnginePhysicalPlan.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/EnginePhysicalPlan.java @@ -22,8 +22,6 @@ import java.util.List; import java.util.stream.Collectors; import java.util.stream.Stream; -import org.apache.flink.configuration.Configuration; -import org.apache.flink.configuration.ConfigurationUtils; /** A jackson serializable object */ public interface EnginePhysicalPlan { @@ -67,10 +65,6 @@ public static String toSqlString(Stream statements) { .filter(statement -> !statement.isBlank()) .collect(Collectors.joining(SQL_STATEMENT_DELIMITER)); } - - public static String toYamlString(Configuration config) { - return String.join("\n", ConfigurationUtils.convertConfigToWritableLines(config, false)); - } } enum ArtifactType { diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/PhysicalPlan.java b/sqrl-planner/src/main/java/com/datasqrl/engine/PhysicalPlan.java index bf44927935..3031e9bc41 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/PhysicalPlan.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/PhysicalPlan.java @@ -16,6 +16,7 @@ package com.datasqrl.engine; import com.datasqrl.canonicalizer.Name; +import com.datasqrl.deployment.model.MutationDatabaseModel; import com.datasqrl.engine.pipeline.ExecutionStage; import com.datasqrl.plan.global.PhysicalPlanRewriter; import com.datasqrl.planner.Sqrl2FlinkSQLTranslator; @@ -42,7 +43,7 @@ public Stream getPlans(Class clazz) { return StreamUtil.filterByClass(stagePlans.stream().map(PhysicalStagePlan::plan), clazz); } - public MutationDatabase getMutationDatabase() { + public MutationDatabaseModel getMutationDatabase() { return MutationDatabase.from(mutationTables.values()); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/EngineCreateTable.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/EngineCreateTable.java index 90ffb3592c..682467cecd 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/EngineCreateTable.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/EngineCreateTable.java @@ -15,7 +15,7 @@ */ package com.datasqrl.engine.database; -/** Used by {@link DatabaseEngine} to keep track of information on created tables */ +/** Used by {@link DatabaseEngine} to keep track of information on created tables. */ public interface EngineCreateTable { EngineCreateTable NONE = new EngineCreateTable() {}; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCEngine.java index 8ca7bd6aff..b7566965bc 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCEngine.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJDBCEngine.java @@ -140,7 +140,8 @@ public EnginePhysicalPlan plan(MaterializationStagePlan stagePlan) { planBuilder.statement(stmt); if (stmt instanceof CreateTableJdbcStatement createTbl) { - tableIdMap.put(createTbl.getEngineTable().table().getTableName(), createTbl); + var engineTable = createTbl.getEngineTable(); + tableIdMap.put(engineTable.table().getTableName(), createTbl); } if (!tableNames.add(stmt.getName())) { diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java index cdcd374abd..883a2465d8 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java @@ -23,10 +23,10 @@ import com.datasqrl.calcite.convert.SqlNodeToString; import com.datasqrl.calcite.dialect.postgres.SqlCreatePostgresView; import com.datasqrl.canonicalizer.Name; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.CreateTableDdlFactory; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; -import com.datasqrl.engine.database.relational.JdbcStatement.Field; -import com.datasqrl.engine.database.relational.JdbcStatement.Type; import com.datasqrl.engine.database.relational.ddl.GenericCreateViewDdlFactory; import com.datasqrl.planner.dag.plan.MaterializationStagePlan.Query; import com.datasqrl.planner.hint.DataTypeHint; @@ -120,7 +120,6 @@ public QueryResult createPassthroughQuery(Query query, boolean withView) { Type.VIEW, viewSql, description, - rowType, getColumns( rowType.getFieldList(), PlannerHints.EMPTY, query.function().getDocumentation())); @@ -196,8 +195,7 @@ protected List getColumns( .collect(Collectors.toList()); } - protected JdbcStatement.Field toField( - RelDataTypeField field, PlannerHints hints, Documentation documentation) { + protected Field toField(RelDataTypeField field, PlannerHints hints, Documentation documentation) { var castSpec = getSqlType( field.getType(), @@ -315,7 +313,6 @@ private JdbcStatement getViewStatement( Type.VIEW, viewSql, documentation.getDocString(null), - rowType, getColumns(rowType.getFieldList(), PlannerHints.EMPTY, documentation)); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/CreateTableJdbcStatement.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/CreateTableJdbcStatement.java index d216243ea5..772a01b74a 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/CreateTableJdbcStatement.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/CreateTableJdbcStatement.java @@ -15,6 +15,9 @@ */ package com.datasqrl.engine.database.relational; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; @@ -103,13 +106,6 @@ public String getSql(CreateTableDdlFactory ddlFactory) { return ddlFactory.createTableDdl(this); } - public enum PartitionType { - NONE, - HASH, - LIST, - RANGE - } - public interface CreateTableDdlFactory { String createTableDdl(CreateTableJdbcStatement stmt); diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/GenericJdbcStatement.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/GenericJdbcStatement.java index 4b286be3c8..ded4d8fad0 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/GenericJdbcStatement.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/GenericJdbcStatement.java @@ -15,6 +15,8 @@ */ package com.datasqrl.engine.database.relational; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlan.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlan.java index 879548550b..4523dfc773 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlan.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlan.java @@ -15,12 +15,12 @@ */ package com.datasqrl.engine.database.relational; +import com.datasqrl.deployment.model.JdbcPlanModel; +import com.datasqrl.deployment.model.JdbcStatementModel; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.database.DatabasePhysicalPlan; -import com.datasqrl.engine.database.relational.JdbcStatement.Type; import com.datasqrl.engine.pipeline.ExecutionStage; -import com.fasterxml.jackson.annotation.JsonCreator; -import com.fasterxml.jackson.annotation.JsonIgnore; -import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonValue; import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -38,34 +38,33 @@ */ @Builder(toBuilder = true) public record JdbcPhysicalPlan( - @JsonIgnore ExecutionStage stage, + ExecutionStage stage, @Singular List statements, @Singular List standaloneExtensionStatements, - @JsonIgnore @Singular List queries, - @JsonIgnore Map tableIdMap) + @Singular List queries, + Map tableIdMap) implements DatabasePhysicalPlan { - @SuppressWarnings("unused") - @JsonCreator - public JdbcPhysicalPlan(@JsonProperty("statements") List statements) { - this(null, new ArrayList<>(statements), List.of(), List.of(), Map.of()); + @JsonValue + public JdbcPlanModel toModel() { + return new JdbcPlanModel( + statements.stream().map(JdbcPhysicalPlan::toStatementModel).toList(), + standaloneExtensionStatements.stream().map(JdbcPhysicalPlan::toStatementModel).toList()); } public List getStatementsForType(Type type) { - return statements.stream().filter(s -> s.getType() == type).collect(Collectors.toList()); + return statements.stream().filter(statement -> statement.getType() == type).toList(); } - @JsonIgnore - @Override public List getDeploymentArtifacts() { var artifacts = new ArrayList(); artifacts.add(new DeploymentArtifact("-schema.sql", buildSchemaContent())); artifacts.add(new DeploymentArtifact("-views.sql", toSql(getStatementsForType(Type.VIEW)))); - standaloneExtensionStatements.stream() - .map(stmt -> new DeploymentArtifact(formatSuffix(stmt.getName()), toSql(stmt))) + .map( + statement -> + new DeploymentArtifact("-" + statement.getName() + ".sql", toSql(statement))) .forEach(artifacts::add); - return List.copyOf(artifacts); } @@ -77,15 +76,42 @@ private String buildSchemaContent() { .collect(Collectors.joining(";\n\n")); } - private static String toSql(List statements) { - return DeploymentArtifact.toSqlString(statements.stream().map(JdbcStatement::getSql)); + private static JdbcStatementModel toStatementModel(JdbcStatement statement) { + var fields = + statement.getFields() == null + ? null + : statement.getFields().stream() + .map( + field -> + new JdbcStatementModel.Field( + field.name(), field.type(), field.nullable(), field.description())) + .toList(); + if (statement instanceof CreateTableJdbcStatement createTable) { + return new JdbcStatementModel( + statement.getName(), + statement.getType(), + statement.getSql(), + statement.getDescription(), + fields, + createTable.getPrimaryKey(), + createTable.getPartitionKey(), + createTable.getPartitionType(), + createTable.getNumPartitions(), + createTable.getTtl()); + } + return new JdbcStatementModel( + statement.getName(), + statement.getType(), + statement.getSql(), + statement.getDescription(), + fields); } - private static String toSql(JdbcStatement stmt) { - return DeploymentArtifact.toSqlString(stmt.getSql()); + private static String toSql(List statements) { + return DeploymentArtifact.toSqlString(statements.stream().map(JdbcStatement::getSql)); } - private static String formatSuffix(String name) { - return "-" + name + ".sql"; + private static String toSql(JdbcStatement statement) { + return DeploymentArtifact.toSqlString(statement.getSql()); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatement.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatement.java index 5f91a894d3..a927f3e0d3 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatement.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/JdbcStatement.java @@ -15,6 +15,8 @@ */ package com.datasqrl.engine.database.relational; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.fasterxml.jackson.databind.annotation.JsonDeserialize; import java.util.List; @@ -30,14 +32,4 @@ public interface JdbcStatement { String getDescription(); List getFields(); - - enum Type { - TABLE, - VIEW, - QUERY, - INDEX, - EXTENSION - } - - record Field(String name, String type, boolean nullable, String description) {} } 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 dc464c4a85..82885ec9e4 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 @@ -23,8 +23,8 @@ import com.datasqrl.calcite.convert.PostgresSqlNodeToString; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; import com.datasqrl.config.JdbcDialect; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; -import com.datasqrl.engine.database.relational.JdbcStatement.Type; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.database.relational.ddl.CreateIndexDDL; import com.datasqrl.engine.database.relational.ddl.InsertStatement; import com.datasqrl.engine.database.relational.ddl.PostgresCreateTableDdlFactory; 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 6a582391b8..9c9cb33723 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 @@ -23,7 +23,7 @@ import com.datasqrl.calcite.dialect.snowflake.SqlCreateIcebergTableFromObjectStorage; import com.datasqrl.config.JdbcDialect; import com.datasqrl.config.PackageJson.EngineConfig; -import com.datasqrl.engine.database.relational.JdbcStatement.Type; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.database.relational.ddl.GenericCreateTableDdlFactory; import com.datasqrl.plan.global.IndexDefinition; import com.datasqrl.planner.hint.DataTypeHint; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/GenericCreateTableDdlFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/GenericCreateTableDdlFactory.java index c68410e7f3..03cbc89020 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/GenericCreateTableDdlFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/GenericCreateTableDdlFactory.java @@ -18,9 +18,9 @@ import static com.datasqrl.engine.database.relational.AbstractJdbcStatementFactory.quoteIdentifier; import static com.datasqrl.engine.database.relational.AbstractJdbcStatementFactory.quoteIdentifiers; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.CreateTableDdlFactory; -import com.datasqrl.engine.database.relational.JdbcStatement.Field; import java.util.List; import java.util.StringJoiner; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/IcebergCreateTableDdlFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/IcebergCreateTableDdlFactory.java index 4181e8a32e..20b930a4a4 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/IcebergCreateTableDdlFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/IcebergCreateTableDdlFactory.java @@ -15,8 +15,8 @@ */ package com.datasqrl.engine.database.relational.ddl; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; import java.util.EnumSet; import java.util.Set; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactory.java index 52e7488bc0..e49cf24ca4 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactory.java @@ -17,8 +17,8 @@ import static com.datasqrl.engine.database.relational.AbstractJdbcStatementFactory.quoteIdentifier; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; import java.time.Duration; import java.util.EnumSet; import java.util.Optional; diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaLogEngine.java b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaLogEngine.java index a46d9af229..be4df1bf5f 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaLogEngine.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaLogEngine.java @@ -24,6 +24,7 @@ import com.datasqrl.config.TestRunnerConfiguration; import com.datasqrl.datatype.DataTypeMapping; import com.datasqrl.datatype.flink.json.FlexibleJsonFlinkFormatTypeMapper; +import com.datasqrl.deployment.model.KafkaNewTopicModel; import com.datasqrl.engine.EngineFeature; import com.datasqrl.engine.EnginePhysicalPlan; import com.datasqrl.engine.ExecutionEngine; @@ -339,21 +340,21 @@ public EnginePhysicalPlan plan(MaterializationStagePlan stagePlan) { Streams.concat( stagePlan.getTables().stream() .map(Table.class::cast) - .map(t -> createNewTopic(t, NewTopic.Type.SUBSCRIPTION)), + .map(t -> createNewTopic(t, KafkaNewTopicModel.Type.SUBSCRIPTION)), stagePlan.getMutations().stream() .map(Table.class::cast) - .map(t -> createNewTopic(t, NewTopic.Type.MUTATION))) + .map(t -> createNewTopic(t, KafkaNewTopicModel.Type.MUTATION))) .toList(); var testRunnerTopics = testRunnerConfig.getCreateTopics().stream() - .map(topicName -> new NewTopic(topicName, topicName)) + .map(topicName -> new KafkaNewTopic(new KafkaNewTopicModel(topicName, topicName))) .toList(); return new KafkaPhysicalPlan(topics, testRunnerTopics); } - private NewTopic createNewTopic(Table table, NewTopic.Type type) { + private KafkaNewTopic createNewTopic(Table table, KafkaNewTopicModel.Type type) { String messageSchema; if (format.startsWith("avro")) { messageSchema = @@ -361,16 +362,17 @@ private NewTopic createNewTopic(Table table, NewTopic.Type type) { } else { messageSchema = ""; // TODO: generate JSON schema } - return new NewTopic( - table.topicName(), - table.tableName(), - table.format(), - numPartitions, - replicationFactor, - type, - table.messageKeys(), - messageSchema, - table.config()); + return new KafkaNewTopic( + new KafkaNewTopicModel( + table.topicName(), + table.tableName(), + table.format(), + numPartitions, + replicationFactor, + type, + table.messageKeys(), + messageSchema, + table.config())); } public record Table( diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaNewTopic.java b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaNewTopic.java new file mode 100644 index 0000000000..fc8ef305b0 --- /dev/null +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaNewTopic.java @@ -0,0 +1,22 @@ +/* + * 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.log.kafka; + +import com.datasqrl.deployment.model.KafkaNewTopicModel; +import com.datasqrl.engine.database.EngineCreateTable; + +/** Planner lifecycle wrapper for a Kafka topic deployment model. */ +public record KafkaNewTopic(KafkaNewTopicModel topic) implements EngineCreateTable {} diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaPhysicalPlan.java b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaPhysicalPlan.java index 940f0d5bb9..a332f717e8 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaPhysicalPlan.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/log/kafka/KafkaPhysicalPlan.java @@ -15,15 +15,26 @@ */ package com.datasqrl.engine.log.kafka; +import com.datasqrl.deployment.model.KafkaPlanModel; import com.datasqrl.engine.EnginePhysicalPlan; -import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonValue; import java.util.List; +import lombok.Builder; +import lombok.Singular; -public record KafkaPhysicalPlan(List topics, List testRunnerTopics) +@Builder +public record KafkaPhysicalPlan( + @Singular List topics, @Singular List testRunnerTopics) implements EnginePhysicalPlan { - @JsonIgnore public boolean isEmpty() { return topics.isEmpty() && testRunnerTopics.isEmpty(); } + + @JsonValue + public KafkaPlanModel toModel() { + return new KafkaPlanModel( + topics.stream().map(KafkaNewTopic::topic).toList(), + testRunnerTopics.stream().map(KafkaNewTopic::topic).toList()); + } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtension.java b/sqrl-planner/src/main/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtension.java index 64e1162a07..97415aa38d 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtension.java +++ b/sqrl-planner/src/main/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtension.java @@ -15,8 +15,8 @@ */ package com.datasqrl.function.translation.postgres.extensions; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; import com.datasqrl.sql.DatabaseTableExtension; import com.google.auto.service.AutoService; import java.time.Duration; diff --git a/sqrl-planner/src/main/java/com/datasqrl/plan/MainScript.java b/sqrl-planner/src/main/java/com/datasqrl/plan/MainScript.java index c577b66322..277966bf29 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/plan/MainScript.java +++ b/sqrl-planner/src/main/java/com/datasqrl/plan/MainScript.java @@ -15,7 +15,7 @@ */ package com.datasqrl.plan; -import com.datasqrl.planner.dag.plan.MutationDatabase; +import com.datasqrl.deployment.model.MutationDatabaseModel; import java.nio.file.Path; import java.util.Optional; @@ -25,7 +25,7 @@ public interface MainScript { String getContent(); - default Optional getMutationDatabase() { + default Optional getMutationDatabase() { return Optional.empty(); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/plan/MainScriptImpl.java b/sqrl-planner/src/main/java/com/datasqrl/plan/MainScriptImpl.java index d5d169cea0..ac7151d17e 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/plan/MainScriptImpl.java +++ b/sqrl-planner/src/main/java/com/datasqrl/plan/MainScriptImpl.java @@ -16,9 +16,9 @@ package com.datasqrl.plan; import com.datasqrl.config.PackageJson; +import com.datasqrl.deployment.model.MutationDatabaseModel; import com.datasqrl.error.ErrorCollector; import com.datasqrl.loaders.resolver.ResourceResolver; -import com.datasqrl.planner.dag.plan.MutationDatabase; import com.datasqrl.util.ConfigLoaderUtils; import com.datasqrl.util.FileUtil; import java.nio.file.Path; @@ -55,7 +55,7 @@ public Optional getPath() { } @Override - public Optional getMutationDatabase() { + public Optional getMutationDatabase() { return config .getScriptConfig() .getDatabase() diff --git a/sqrl-planner/src/main/java/com/datasqrl/planner/FlinkPhysicalPlan.java b/sqrl-planner/src/main/java/com/datasqrl/planner/FlinkPhysicalPlan.java index 1e9255e903..db35b01e60 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/planner/FlinkPhysicalPlan.java +++ b/sqrl-planner/src/main/java/com/datasqrl/planner/FlinkPhysicalPlan.java @@ -17,12 +17,13 @@ import static com.google.common.base.Preconditions.checkArgument; +import com.datasqrl.deployment.model.FlinkPlanModel; import com.datasqrl.engine.EnginePhysicalPlan; import com.datasqrl.engine.database.relational.IcebergEngineFactory; import com.datasqrl.engine.stream.flink.sql.RelToFlinkSql; import com.datasqrl.planner.tables.FlinkConnectorConfigWrapper; import com.datasqrl.planner.util.CompiledPlanCondenser; -import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonValue; import com.google.common.collect.ArrayListMultimap; import com.google.common.collect.ImmutableList; import com.google.common.collect.ListMultimap; @@ -38,6 +39,7 @@ import org.apache.calcite.sql.SqlNode; import org.apache.calcite.sql.parser.SqlParserPos; import org.apache.flink.configuration.Configuration; +import org.apache.flink.configuration.ConfigurationUtils; import org.apache.flink.sql.parser.ddl.SqlCreateFunction; import org.apache.flink.sql.parser.ddl.table.SqlCreateTable; import org.apache.flink.sql.parser.dml.RichSqlInsert; @@ -66,18 +68,27 @@ public class FlinkPhysicalPlan implements EnginePhysicalPlan { Set connectors; Set formats; Set functions; - @JsonIgnore Optional compiledPlan; - @JsonIgnore Optional explainedPlan; - @JsonIgnore List flinkSqlNoFunctions; - @JsonIgnore Configuration config; - @JsonIgnore ListMultimap flinkSqlBatched; - @JsonIgnore ListMultimap flinkSqlNoFunctionsBatched; + Optional compiledPlan; + Optional explainedPlan; + List flinkSqlNoFunctions; + Configuration config; + ListMultimap flinkSqlBatched; + ListMultimap flinkSqlNoFunctionsBatched; + + @JsonValue + public FlinkPlanModel toModel() { + return new FlinkPlanModel(flinkSql, connectors, formats, functions); + } + + private static String toYamlString(Configuration config) { + return String.join("\n", ConfigurationUtils.convertConfigToWritableLines(config, false)); + } @Override public List getDeploymentArtifacts() { var builder = ImmutableList.builder(); builder.add( - new DeploymentArtifact("-config.yaml", DeploymentArtifact.toYamlString(config)), + new DeploymentArtifact("-config.yaml", toYamlString(config)), new DeploymentArtifact("-sql.sql", DeploymentArtifact.toSqlString(flinkSql)), new DeploymentArtifact( "-sql-no-functions.sql", DeploymentArtifact.toSqlString(flinkSqlNoFunctions)), diff --git a/sqrl-planner/src/main/java/com/datasqrl/planner/dag/plan/MutationDatabase.java b/sqrl-planner/src/main/java/com/datasqrl/planner/dag/plan/MutationDatabase.java index 37470dee85..6c957283f5 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/planner/dag/plan/MutationDatabase.java +++ b/sqrl-planner/src/main/java/com/datasqrl/planner/dag/plan/MutationDatabase.java @@ -16,23 +16,28 @@ package com.datasqrl.planner.dag.plan; import com.datasqrl.calcite.type.TypeCompatibility; +import com.datasqrl.deployment.model.MutationDatabaseModel; +import com.datasqrl.deployment.model.MutationDatabaseModel.ColumnDefinition; +import com.datasqrl.deployment.model.MutationDatabaseModel.Table; +import com.datasqrl.deployment.model.MutationDatabaseModel.TableDefinition; import com.datasqrl.error.ErrorCollector; import com.datasqrl.planner.Sqrl2FlinkSQLTranslator; import com.datasqrl.planner.Sqrl2FlinkSQLTranslator.ParsedRelDataTypeResult; import com.datasqrl.server.exec.FlinkExecFunction; -import com.fasterxml.jackson.annotation.JsonInclude; import java.util.Collection; import java.util.List; -import java.util.Map; import java.util.Objects; import java.util.function.Function; import java.util.stream.Collectors; +import lombok.AccessLevel; +import lombok.NoArgsConstructor; import org.apache.flink.sql.parser.ddl.SqlTableColumn; -@JsonInclude(JsonInclude.Include.NON_EMPTY) -public record MutationDatabase(List
tables) { +/** Builds and compares {@link MutationDatabaseModel}s during planning. */ +@NoArgsConstructor(access = AccessLevel.PRIVATE) +public final class MutationDatabase { - public static MutationDatabase from(Collection mutationTables) { + public static MutationDatabaseModel from(Collection mutationTables) { var tables = mutationTables.stream() .map( @@ -66,16 +71,20 @@ public static MutationDatabase from(Collection mutationTables) { mutTbl.getDocumentation().getDocString(null)); }) .toList(); - return new MutationDatabase(tables); + + return new MutationDatabaseModel(tables); } - public boolean isBackwardsCompatible( - MutationDatabase compareDb, Sqrl2FlinkSQLTranslator env, ErrorCollector errors) { + public static boolean isBackwardsCompatible( + MutationDatabaseModel database, + MutationDatabaseModel compareDb, + Sqrl2FlinkSQLTranslator env, + ErrorCollector errors) { var compareTablesByName = compareDb.tables().stream().collect(Collectors.toMap(Table::canonicalName, t -> t)); var compatible = true; - for (var table : tables) { + for (var table : database.tables()) { var compareTable = compareTablesByName.get(table.canonicalName()); if (compareTable == null) { continue; @@ -160,17 +169,4 @@ public boolean isBackwardsCompatible( return compatible; } - - public record Table( - String canonicalName, - String engine, - String createTableSql, - TableDefinition definition, - Map configOptions, - String documentation) {} - - public record TableDefinition( - List columns, List primaryKey, List partitionKey) {} - - public record ColumnDefinition(String name, String spec, String documentation) {} } diff --git a/sqrl-planner/src/main/java/com/datasqrl/util/ConfigLoaderUtils.java b/sqrl-planner/src/main/java/com/datasqrl/util/ConfigLoaderUtils.java index fe13deb5a9..e87a0f0753 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/util/ConfigLoaderUtils.java +++ b/sqrl-planner/src/main/java/com/datasqrl/util/ConfigLoaderUtils.java @@ -20,15 +20,19 @@ import com.datasqrl.config.PackageJson; import com.datasqrl.config.SqrlConfig; import com.datasqrl.config.SqrlConstants; -import com.datasqrl.engine.database.relational.JdbcPhysicalPlan; -import com.datasqrl.engine.log.kafka.KafkaPhysicalPlan; +import com.datasqrl.deployment.model.JdbcPlanModel; +import com.datasqrl.deployment.model.KafkaPlanModel; +import com.datasqrl.deployment.model.MutationDatabaseModel; import com.datasqrl.error.ErrorCollector; import com.datasqrl.error.ErrorMessage; import com.datasqrl.error.ResourceFileUtil; -import com.datasqrl.planner.dag.plan.MutationDatabase; import com.datasqrl.planner.util.NonSecretEnvVarResolver; +import com.fasterxml.jackson.annotation.JsonAutoDetect.Visibility; +import com.fasterxml.jackson.annotation.PropertyAccessor; +import com.fasterxml.jackson.databind.DeserializationFeature; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.annotations.VisibleForTesting; @@ -57,15 +61,11 @@ public final class ConfigLoaderUtils { public static final String PACKAGE_SCHEMA_PATH = "/jsonSchema/packageSchema.json"; public static final ObjectMapper MAPPER = - new ObjectMapper() - .setVisibility( - com.fasterxml.jackson.annotation.PropertyAccessor.FIELD, - com.fasterxml.jackson.annotation.JsonAutoDetect.Visibility.ANY) - .configure( - com.fasterxml.jackson.databind.DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, - false) - .configure( - com.fasterxml.jackson.databind.SerializationFeature.FAIL_ON_EMPTY_BEANS, false); + JsonUtils.MAPPER + .copy() + .setVisibility(PropertyAccessor.FIELD, Visibility.ANY) + .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false) + .configure(SerializationFeature.FAIL_ON_EMPTY_BEANS, false); private static final List DEFAULTS = List.of("/default-package.json"); private static final List RUN_DEFAULTS = @@ -216,24 +216,24 @@ public static Configuration loadFlinkConfig(Path planDir) throws IOException { } /** - * Loads the Kafka physical plan from the plan directory by parsing the {@code kafka.json} - * configuration file. This method constructs a complete {@link KafkaPhysicalPlan} containing both + * Loads the Kafka deployment file from the plan directory by parsing the {@code kafka.json} + * configuration file. This method constructs a complete {@link KafkaPlanModel} containing both * regular topics and test runner topics. * * @param planDir the plan directory containing the {@code kafka.json} file - * @return an {@link Optional} containing a {@link KafkaPhysicalPlan} with topics and test runner + * @return an {@link Optional} containing a {@link KafkaPlanModel} with topics and test runner * topics, or empty if no {@code kafka.json} file exists * @throws IllegalArgumentException if the plan directory is null or does not exist * @throws IllegalStateException if the {@code kafka.json} file exists but cannot be parsed */ - public static Optional loadKafkaPhysicalPlan(Path planDir) { + public static Optional loadKafkaPhysicalPlan(Path planDir) { validatePlanDir(planDir); var kafkaFile = planDir.resolve("kafka.json").toFile(); if (kafkaFile.exists()) { try { - var kafkaPlan = MAPPER.readValue(kafkaFile, KafkaPhysicalPlan.class); + var kafkaPlan = MAPPER.readValue(kafkaFile, KafkaPlanModel.class); return Optional.of(kafkaPlan); } catch (Exception ex) { @@ -244,14 +244,15 @@ public static Optional loadKafkaPhysicalPlan(Path planDir) { return Optional.empty(); } - public static MutationDatabase loadMutationDatabase(Path databaseFile, ErrorCollector errors) { + public static MutationDatabaseModel loadMutationDatabase( + Path databaseFile, ErrorCollector errors) { checkArgument( Files.isRegularFile(databaseFile), "Mutation database file does not exist: %s", databaseFile); try { - return MAPPER.readValue(databaseFile.toFile(), MutationDatabase.class); + return MAPPER.readValue(databaseFile.toFile(), MutationDatabaseModel.class); } catch (IOException e) { throw errors.exception( "Failed to load mutation database from '%s': %s", databaseFile, e.getMessage()); @@ -259,17 +260,18 @@ public static MutationDatabase loadMutationDatabase(Path databaseFile, ErrorColl } /** - * Loads PostgreSQL physical plan from the plan directory by parsing the {@code postgres.json} - * configuration file. This method constructs a complete {@link JdbcPhysicalPlan} containing - * database statements for schema creation, views, indexes, and other database artifacts. + * Loads the PostgreSQL deployment file from the plan directory by parsing the {@code + * postgres.json} configuration file. This method constructs a complete {@link JdbcPlanModel} + * containing database statements for schema creation, views, indexes, and other database + * artifacts. * * @param planDir the plan directory containing the {@code postgres.json} file - * @return an {@link Optional} containing a {@link JdbcPhysicalPlan} with database statements, or + * @return an {@link Optional} containing a {@link JdbcPlanModel} with database statements, or * empty if no {@code postgres.json} file exists * @throws IllegalArgumentException if the plan directory is null or does not exist * @throws IllegalStateException if the {@code postgres.json} file exists but cannot be parsed */ - public static Optional loadPostgresPhysicalPlan(Path planDir) { + public static Optional loadPostgresPhysicalPlan(Path planDir) { validatePlanDir(planDir); var postgresFile = planDir.resolve("postgres.json").toFile(); @@ -278,7 +280,7 @@ public static Optional loadPostgresPhysicalPlan(Path planDir) } try { - var jdbcPlan = MAPPER.readValue(postgresFile, JdbcPhysicalPlan.class); + var jdbcPlan = MAPPER.readValue(postgresFile, JdbcPlanModel.class); return Optional.of(jdbcPlan); } catch (Exception ex) { diff --git a/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlanTest.java b/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlanTest.java index 9d8dea92e7..61416c63de 100644 --- a/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlanTest.java +++ b/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/JdbcPhysicalPlanTest.java @@ -17,8 +17,8 @@ import static org.assertj.core.api.Assertions.assertThat; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; import com.datasqrl.engine.EnginePhysicalPlan.DeploymentArtifact; -import com.datasqrl.engine.database.relational.JdbcStatement.Type; import java.util.List; import java.util.Map; import org.junit.jupiter.api.Test; @@ -74,7 +74,7 @@ void givenStandaloneExtensionStatements_whenGetDeploymentArtifacts_thenAppendedA void givenJsonCreatorPlan_whenGetStatementsForType_thenFiltersByType() { var table = stmt("t1", Type.TABLE, "CREATE TABLE t1"); var view = stmt("v1", Type.VIEW, "CREATE VIEW v1"); - var plan = new JdbcPhysicalPlan(List.of(table, view)); + var plan = new JdbcPhysicalPlan(null, List.of(table, view), List.of(), List.of(), Map.of()); assertThat(plan.getStatementsForType(Type.TABLE)).containsExactly(table); assertThat(plan.getStatementsForType(Type.INDEX)).isEmpty(); diff --git a/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactoryTest.java b/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactoryTest.java index f7b7377d51..01ce088f73 100644 --- a/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactoryTest.java +++ b/sqrl-planner/src/test/java/com/datasqrl/engine/database/relational/ddl/PostgresCreateTableDdlFactoryTest.java @@ -18,9 +18,9 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; -import com.datasqrl.engine.database.relational.JdbcStatement.Field; import java.time.Duration; import java.util.List; import org.junit.jupiter.api.Test; diff --git a/sqrl-planner/src/test/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtensionTest.java b/sqrl-planner/src/test/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtensionTest.java index 2f8b5c20cb..40475bb7ba 100644 --- a/sqrl-planner/src/test/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtensionTest.java +++ b/sqrl-planner/src/test/java/com/datasqrl/function/translation/postgres/extensions/PgPartmanExtensionTest.java @@ -17,8 +17,8 @@ import static org.assertj.core.api.Assertions.assertThat; +import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; -import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; import java.time.Duration; import java.util.List; import org.junit.jupiter.api.Test; diff --git a/sqrl-planner/src/test/java/com/datasqrl/planner/PlanModelSerializationTest.java b/sqrl-planner/src/test/java/com/datasqrl/planner/PlanModelSerializationTest.java new file mode 100644 index 0000000000..e52867f42d --- /dev/null +++ b/sqrl-planner/src/test/java/com/datasqrl/planner/PlanModelSerializationTest.java @@ -0,0 +1,161 @@ +/* + * 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.planner; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.datasqrl.deployment.model.FlinkPlanModel; +import com.datasqrl.deployment.model.JdbcPlanModel; +import com.datasqrl.deployment.model.JdbcStatementModel; +import com.datasqrl.deployment.model.JdbcStatementModel.Field; +import com.datasqrl.deployment.model.JdbcStatementModel.Type; +import com.datasqrl.deployment.model.KafkaNewTopicModel; +import com.datasqrl.engine.database.relational.CreateTableJdbcStatement; +import com.datasqrl.engine.database.relational.GenericJdbcStatement; +import com.datasqrl.engine.database.relational.JdbcPhysicalPlan; +import com.datasqrl.engine.log.kafka.KafkaNewTopic; +import com.datasqrl.engine.log.kafka.KafkaPhysicalPlan; +import com.datasqrl.util.SqrlObjectMapper; +import java.time.Duration; +import java.util.List; +import java.util.Map; +import java.util.Set; +import org.junit.jupiter.api.Test; + +/** Guards the explicit mappings from planner state to deployment-file models. */ +class PlanModelSerializationTest { + + @Test + void givenJdbcPhysicalPlan_whenMapped_thenReturnsJdbcPlanModel() { + var field = new Field("id", "BIGINT", false, "primary key"); + var createTable = + new CreateTableJdbcStatement( + "orders", + "order storage", + List.of(field), + List.of("id"), + List.of("tenant"), + JdbcStatementModel.PartitionType.HASH, + 4, + Duration.ofSeconds(30), + null, + statement -> "CREATE TABLE orders"); + var view = + new GenericJdbcStatement( + "orders_view", + Type.VIEW, + "CREATE VIEW orders_view", + "order projection", + List.of(field)); + var extension = + new GenericJdbcStatement( + "pg_cron", Type.EXTENSION, "CREATE EXTENSION pg_cron", "scheduler", List.of()); + var plan = + JdbcPhysicalPlan.builder() + .statement(createTable) + .statement(view) + .standaloneExtensionStatement(extension) + .tableIdMap(Map.of()) + .build(); + + JdbcPlanModel model = plan.toModel(); + + assertThat(model.statements()).hasSize(2); + var table = model.statements().get(0); + assertThat(table.name()).isEqualTo("orders"); + assertThat(table.type()).isEqualTo(JdbcStatementModel.Type.TABLE); + assertThat(table.sql()).isEqualTo("CREATE TABLE orders"); + assertThat(table.description()).isEqualTo("order storage"); + assertThat(table.fields()) + .containsExactly(new JdbcStatementModel.Field("id", "BIGINT", false, "primary key")); + assertThat(table.primaryKey()).containsExactly("id"); + assertThat(table.partitionKey()).containsExactly("tenant"); + assertThat(table.partitionType()).isEqualTo(JdbcStatementModel.PartitionType.HASH); + assertThat(table.numPartitions()).isEqualTo(4); + assertThat(table.ttl()).isEqualTo(Duration.ofSeconds(30)); + + var viewModel = model.statements().get(1); + assertThat(viewModel.name()).isEqualTo("orders_view"); + assertThat(viewModel.type()).isEqualTo(JdbcStatementModel.Type.VIEW); + assertThat(viewModel.sql()).isEqualTo("CREATE VIEW orders_view"); + assertThat(viewModel.description()).isEqualTo("order projection"); + assertThat(viewModel.fields()) + .containsExactly(new JdbcStatementModel.Field("id", "BIGINT", false, "primary key")); + assertThat(viewModel.primaryKey()).isNull(); + assertThat(viewModel.partitionKey()).isNull(); + assertThat(viewModel.partitionType()).isNull(); + assertThat(viewModel.numPartitions()).isNull(); + assertThat(viewModel.ttl()).isNull(); + + assertThat(model.standaloneExtensionStatements()) + .containsExactly( + new JdbcStatementModel( + "pg_cron", + JdbcStatementModel.Type.EXTENSION, + "CREATE EXTENSION pg_cron", + "scheduler", + List.of())); + } + + @Test + void givenFlinkPhysicalPlan_whenMapped_thenReturnsFlinkPlanModel() { + var plan = + FlinkPhysicalPlan.builder() + .flinkSql(List.of("CREATE TABLE x")) + .connectors(Set.of("kafka")) + .formats(Set.of("json")) + .functions(Set.of("CREATE FUNCTION f")) + .build(); + + FlinkPlanModel model = plan.toModel(); + + assertThat(model.flinkSql()).containsExactly("CREATE TABLE x"); + assertThat(model.connectors()).containsExactly("kafka"); + assertThat(model.formats()).containsExactly("json"); + assertThat(model.functions()).containsExactly("CREATE FUNCTION f"); + } + + @Test + void givenKafkaPhysicalPlan_whenMapped_thenReturnsWrappedTopicModel() { + var topic = new KafkaNewTopicModel("orders", "orders", 3, (short) 2); + var testRunnerTopic = new KafkaNewTopicModel("test-orders", "test-orders"); + var plan = + KafkaPhysicalPlan.builder() + .topic(new KafkaNewTopic(topic)) + .testRunnerTopic(new KafkaNewTopic(testRunnerTopic)) + .build(); + + var model = plan.toModel(); + + assertThat(model.topics()).containsExactly(topic); + assertThat(model.testRunnerTopics()).containsExactly(testRunnerTopic); + assertThat(topic.format()).isNull(); + assertThat(topic.numPartitions()).isEqualTo(3); + assertThat(topic.replicationFactor()).isEqualTo((short) 2); + assertThat(topic.type()).isEqualTo(KafkaNewTopicModel.Type.SUBSCRIPTION); + assertThat(topic.messageKeys()).isEmpty(); + assertThat(topic.messageSchema()).isEmpty(); + assertThat(topic.config()).isEmpty(); + assertThat(testRunnerTopic.numPartitions()).isEqualTo(1); + assertThat(testRunnerTopic.replicationFactor()).isEqualTo((short) 1); + + var json = SqrlObjectMapper.INSTANCE.valueToTree(model); + + assertThat(List.copyOf(json.properties())) + .extracting(Map.Entry::getKey) + .containsExactly("topics", "testRunnerTopics"); + } +} diff --git a/sqrl-planner/src/test/java/com/datasqrl/util/ConfigLoaderUtilsTest.java b/sqrl-planner/src/test/java/com/datasqrl/util/ConfigLoaderUtilsTest.java index 2080004bd0..0e17b74ce2 100644 --- a/sqrl-planner/src/test/java/com/datasqrl/util/ConfigLoaderUtilsTest.java +++ b/sqrl-planner/src/test/java/com/datasqrl/util/ConfigLoaderUtilsTest.java @@ -22,13 +22,14 @@ import com.datasqrl.config.PackageJsonImpl; import com.datasqrl.config.SqrlConfig; import com.datasqrl.config.SqrlConfigTest; -import com.datasqrl.engine.database.relational.JdbcStatement; -import com.datasqrl.engine.log.kafka.NewTopic; +import com.datasqrl.deployment.model.JdbcStatementModel; +import com.datasqrl.deployment.model.KafkaNewTopicModel; import com.datasqrl.error.CollectedException; import com.datasqrl.error.ErrorCollector; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.time.Duration; import java.util.List; import lombok.SneakyThrows; import org.apache.flink.configuration.Configuration; @@ -324,7 +325,7 @@ void givenPlanDirWithValidKafkaJson_whenLoadKafkaPhysicalPlan_thenReturnsPopulat assertThat(kafkaPlan.topics().get(1).topicName()).isEqualTo("customers"); assertThat(kafkaPlan.topics().get(1).numPartitions()).isEqualTo(1); - assertThat(kafkaPlan.testRunnerTopics().stream().map(NewTopic::topicName)) + assertThat(kafkaPlan.testRunnerTopics().stream().map(KafkaNewTopicModel::topicName)) .containsExactlyInAnyOrder("test-topic-1", "test-topic-2"); assertThat(kafkaPlan.isEmpty()).isFalse(); } @@ -373,7 +374,7 @@ void givenPlanDirWithoutPostgresJson_whenLoadPostgresPhysicalPlan_thenReturnsEmp } @Test - void givenPlanDirWithValidPostgresJson_whenLoadPostgresPhysicalPlan_thenReturnsJdbcPlan() + void givenPlanDirWithValidPostgresJson_whenLoadPostgresPhysicalPlan_thenReturnsJdbcPlanModel() throws IOException { // Given Path planDir = tempDir.resolve("plan"); @@ -384,9 +385,10 @@ void givenPlanDirWithValidPostgresJson_whenLoadPostgresPhysicalPlan_thenReturnsJ { "statements": [ { - "name": "create_users_table", - "type": "TABLE", - "sql": "CREATE TABLE users (id INT PRIMARY KEY, name VARCHAR(100))" + "name": "create_users_table", + "type": "TABLE", + "sql": "CREATE TABLE users (id INT PRIMARY KEY, name VARCHAR(100))", + "ttl": 30.0 }, { "name": "create_orders_table", @@ -414,17 +416,18 @@ void givenPlanDirWithValidPostgresJson_whenLoadPostgresPhysicalPlan_thenReturnsJ assertThat(jdbcPlan.statements()).hasSize(3); var statements = jdbcPlan.statements(); - assertThat(statements.get(0).getName()).isEqualTo("create_users_table"); - assertThat(statements.get(0).getType()).isEqualTo(JdbcStatement.Type.TABLE); - assertThat(statements.get(0).getSql()).contains("CREATE TABLE users"); - - assertThat(statements.get(1).getName()).isEqualTo("create_orders_table"); - assertThat(statements.get(1).getType()).isEqualTo(JdbcStatement.Type.TABLE); - assertThat(statements.get(1).getSql()).contains("CREATE TABLE orders"); - - assertThat(statements.get(2).getName()).isEqualTo("create_index"); - assertThat(statements.get(2).getType()).isEqualTo(JdbcStatement.Type.INDEX); - assertThat(statements.get(2).getSql()).contains("CREATE INDEX"); + assertThat(statements.get(0).name()).isEqualTo("create_users_table"); + assertThat(statements.get(0).type()).isEqualTo(JdbcStatementModel.Type.TABLE); + assertThat(statements.get(0).sql()).contains("CREATE TABLE users"); + assertThat(statements.get(0).ttl()).isEqualTo(Duration.ofSeconds(30)); + + assertThat(statements.get(1).name()).isEqualTo("create_orders_table"); + assertThat(statements.get(1).type()).isEqualTo(JdbcStatementModel.Type.TABLE); + assertThat(statements.get(1).sql()).contains("CREATE TABLE orders"); + + assertThat(statements.get(2).name()).isEqualTo("create_index"); + assertThat(statements.get(2).type()).isEqualTo(JdbcStatementModel.Type.INDEX); + assertThat(statements.get(2).sql()).contains("CREATE INDEX"); } @Test