From 7d17157633da687c881fb3ed8dee0aa518d44c63 Mon Sep 17 00:00:00 2001 From: 0264408 Date: Wed, 16 Sep 2026 11:05:14 +0800 Subject: [PATCH] feat: add a pluggable JDBC dialect layer with MySQL support Fold Flyway V1-V7 into per-dialect baselines and resolve Boot 3/4 Flyway customizer types without a DataSource creation cycle. Co-authored-by: Cursor --- README.md | 10 +- docs/roadmap.md | 6 +- docs/usage.md | 14 +- ordo-spring-boot-autoconfigure/pom.xml | 8 + .../spring/OrdoFlywayAutoConfiguration.java | 123 ++++++++ .../spring/OrdoJdbcAutoConfiguration.java | 36 ++- .../jetlumen/ordo/spring/OrdoProperties.java | 18 ++ ...ot.autoconfigure.AutoConfiguration.imports | 1 + .../spring/OrdoJdbcAutoConfigurationTest.java | 11 +- ordo-spring-boot-starter/pom.xml | 12 +- ordo-storage-jdbc/pom.xml | 5 + .../jdbc/JdbcActionExecutionRepository.java | 12 +- .../jdbc/JdbcApprovalTaskRepository.java | 20 +- .../jdbc/JdbcProcessDefinitionRepository.java | 41 ++- .../jdbc/JdbcProcessHistoryRepository.java | 12 +- .../jdbc/JdbcProcessInstanceRepository.java | 23 +- .../jdbc/dialect/LimitOffsetSqlDialect.java | 19 ++ .../storage/jdbc/dialect/MysqlSqlDialect.java | 38 +++ .../jdbc/dialect/PostgresSqlDialect.java | 52 ++++ .../ordo/storage/jdbc/dialect/SqlDialect.java | 29 ++ .../storage/jdbc/dialect/SqlDialects.java | 78 +++++ ...lumen.ordo.storage.jdbc.dialect.SqlDialect | 2 + .../db/migration/V1__create_ordo_tables.sql | 48 --- .../V2__add_step_candidates_and_policy.sql | 23 -- .../db/migration/V3__add_step_transitions.sql | 12 - .../V4__add_step_kind_and_action_key.sql | 3 - ...add_process_event_and_action_execution.sql | 29 -- .../db/migration/V6__add_step_due.sql | 8 - .../db/migration/V7__definition_versions.sql | 96 ------ .../db/mysql/migration/V1__baseline.sql | 125 ++++++++ .../db/postgresql/migration/V1__baseline.sql | 120 ++++++++ .../jdbc/JdbcMysqlIntegrationTest.java | 282 ++++++++++++++++++ .../jdbc/JdbcPostgresIntegrationTest.java | 2 +- .../ordo/storage/jdbc/JdbcTestSupport.java | 35 +-- .../storage/jdbc/dialect/SqlDialectsTest.java | 68 +++++ pom.xml | 6 + 36 files changed, 1123 insertions(+), 304 deletions(-) create mode 100644 ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoFlywayAutoConfiguration.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/LimitOffsetSqlDialect.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/MysqlSqlDialect.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/PostgresSqlDialect.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialect.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialects.java create mode 100644 ordo-storage-jdbc/src/main/resources/META-INF/services/com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V2__add_step_candidates_and_policy.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V3__add_step_transitions.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V4__add_step_kind_and_action_key.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V6__add_step_due.sql delete mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V7__definition_versions.sql create mode 100644 ordo-storage-jdbc/src/main/resources/db/mysql/migration/V1__baseline.sql create mode 100644 ordo-storage-jdbc/src/main/resources/db/postgresql/migration/V1__baseline.sql create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcMysqlIntegrationTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialectsTest.java diff --git a/README.md b/README.md index c0476f7..b11febf 100644 --- a/README.md +++ b/README.md @@ -10,11 +10,11 @@ |---|---| | `ordo-api` | 公共模型与 `OrdoEngine` 端口 | | `ordo-core` | 运行时(`DefaultOrdoEngine` / `InMemoryOrdoEngine`) | -| `ordo-storage-jdbc` | JDBC 存储 + Flyway 迁移(V1–V7) | +| `ordo-storage-jdbc` | JDBC 存储 + 按方言 Flyway 基线(PostgreSQL / MySQL) | | `ordo-spring-boot-starter` | Spring Boot 自动装配(JDBC + Flyway) | | `ordo-example` | 内存引擎示例 | -存储实现面向 **PostgreSQL**(测试可用 H2 PostgreSQL 兼容模式)。Starter **不携带** JDBC 驱动,宿主自行加入 `postgresql` 或 `h2`。Spring Boot 4 还需额外引入 `spring-boot-starter-flyway`,否则迁移不会跑。 +存储实现通过 dialect 层支持 **PostgreSQL** 与 **MySQL**(测试可用 H2 PostgreSQL 兼容模式)。Starter **不携带** JDBC 驱动,宿主自行加入 `postgresql`、`mysql-connector-j` 或 `h2`。Spring Boot 4 还需额外引入 `spring-boot-starter-flyway`,否则迁移不会跑。 ## 能力 @@ -29,7 +29,7 @@ - 审计时间线:`ProcessEvent` + `queryHistory` - 扩展点:`AssigneeResolver`、`RoutingCondition`、`ActionHandler`、`OrdoEventListener` -开发计划:MySQL 方言、可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST 与目录之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。 +开发计划:可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST 与目录之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。 详细用法(定义 JSON、扩展点、异常、查询、ACTION/审计语义)见 **[docs/usage.md](docs/usage.md)**。对外行为变更时同步更新该文档。 @@ -100,8 +100,10 @@ ordo.approve(manager.id(), "maria", "ok"); ```yaml ordo: enabled: true + jdbc: + dialect: # 可选 postgresql / mysql;空则按 DataSource 探测 definitions: - location: classpath*:ordo/*.json # 默认值;启动时 replace 加载 + location: classpath*:ordo/*.json # 默认值;启动时 publish 加载 ``` 宿主提供 Bean 即可覆盖默认值: diff --git a/docs/roadmap.md b/docs/roadmap.md index c341745..1677d8a 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -7,7 +7,7 @@ ## 已完成 - ANY/ALL 会签/或签 -- JDBC 存储(PostgreSQL)+ Flyway V1–V7 +- JDBC 存储(PostgreSQL / MySQL)+ 按方言 Flyway 基线 - 定义不可变多版本:`publish`;实例锁定 `definitionVersion` - 条件路由 `StepTransition` + `RoutingCondition` - ACTION + `ActionHandler`;执行记录持久化 @@ -22,10 +22,6 @@ 下列能力已纳入计划,尚未实现。实现顺序可按依赖调整,但范围本身不从计划中拿掉。 -### MySQL 方言 - -当前 JDBC DDL/冲突处理面向 PostgreSQL。计划增加 MySQL 方言(及对应测试),通过 dialect 层扩展,而不是只支持一种库。 - ### 可选 REST + 目录 SPI 引擎入口仍是 `OrdoEngine`。计划提供**可选、极薄**的 REST 适配(例如独立 starter),带 OpenAPI,覆盖定义读写/解析校验、实例与任务查询,不包含鉴权、RBAC、业务表单。 diff --git a/docs/usage.md b/docs/usage.md index 90aa18d..dce4567 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -12,7 +12,7 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。 产品边界:不做业务表单、用户体系、多租户;不内置设计器 UI。业务字段放在 `ProcessContext`(不可变 `Map`)。当前也**没有** REST;HTTP 仍由宿主自建。计划中的可选 REST 与独立设计器见 [roadmap.md](roadmap.md)。 -开发计划(尚未提供,见 [roadmap.md](roadmap.md)):MySQL 方言、可选 REST + 目录 SPI。 +开发计划(尚未提供,见 [roadmap.md](roadmap.md)):可选 REST + 目录 SPI。 ## 2. 模块与接入 @@ -20,11 +20,11 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。 |---|---| | `ordo-api` | 始终:模型与 `OrdoEngine` | | `ordo-core` | 内存引擎 / 自己装配 `DefaultOrdoEngine` | -| `ordo-storage-jdbc` | JDBC 持久化(PostgreSQL;测试可用 H2 PostgreSQL 模式) | +| `ordo-storage-jdbc` | JDBC 持久化(PostgreSQL、MySQL;测试可用 H2 PostgreSQL 模式) | | `ordo-spring-boot-starter` | Spring Boot 自动装配 | | `ordo-example` | `LeaveRequestExample` 内存演示 | -Starter **不携带** JDBC 驱动。生产加 `org.postgresql:postgresql`。Spring Boot 4 还需 `spring-boot-starter-flyway`,否则 Flyway 迁移不会执行。 +Starter **不携带** JDBC 驱动。生产按库添加 `org.postgresql:postgresql` 或 `com.mysql:mysql-connector-j`。Spring Boot 4 还需 `spring-boot-starter-flyway`,否则 Flyway 迁移不会执行。 先 `mvn install` 本仓库,宿主再依赖 `0.0.1-SNAPSHOT`。 @@ -60,6 +60,8 @@ OrdoEngine ordo = new InMemoryOrdoEngine(); ```yaml ordo: enabled: true + jdbc: + dialect: # 可选 postgresql / mysql;空则按 DataSource 探测 definitions: location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 publish due: @@ -345,12 +347,14 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的 ## 12. 存储 -Flyway 脚本在 `ordo-storage-jdbc` 的 `db/migration`(V1–V7)。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤/候选人/转移、实例(含 `definition_version`)、任务、`ordo_process_event`、`ordo_action_execution`。 +Flyway 脚本按方言分目录:`db/postgresql/migration`、`db/mysql/migration`(各一份当前 schema 的 `V1__baseline.sql`)。未设置 `spring.flyway.locations` 时,starter 按探测到的方言指向对应目录。已有 `flyway_schema_history` 的开发库需清空后重跑。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤/候选人/转移、实例(含 `definition_version`)、任务、`ordo_process_event`、`ordo_action_execution`。 + +自定义方言:实现 `SqlDialect` 并用 `META-INF/services` 注册,或提供 `SqlDialect` Bean。新增列时每个已支持方言目录各加一条迁移。 多 JVM 共享同一库时,多步写入走 `TransactionExecutor`,完成任务/实例用条件更新(仍 PENDING / 仍 RUNNING 才改),避免双花。 ## 13. 未提供能力 -开发计划中(见 [roadmap.md](roadmap.md)):MySQL 方言、可选 REST + 目录 SPI。独立设计器不进本仓库,等 REST 与目录之后再做。 +开发计划中(见 [roadmap.md](roadmap.md)):可选 REST + 目录 SPI。独立设计器不进本仓库,等 REST 与目录之后再做。 暂不在计划中:多租户、设计器 UI。当前 REST 由宿主自建。 diff --git a/ordo-spring-boot-autoconfigure/pom.xml b/ordo-spring-boot-autoconfigure/pom.xml index c456e66..6170f3c 100644 --- a/ordo-spring-boot-autoconfigure/pom.xml +++ b/ordo-spring-boot-autoconfigure/pom.xml @@ -58,6 +58,14 @@ org.flywaydb flyway-core + + org.flywaydb + flyway-database-postgresql + + + org.flywaydb + flyway-mysql + org.springframework.boot spring-boot-configuration-processor diff --git a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoFlywayAutoConfiguration.java b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoFlywayAutoConfiguration.java new file mode 100644 index 0000000..7ae9532 --- /dev/null +++ b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoFlywayAutoConfiguration.java @@ -0,0 +1,123 @@ +package com.jetlumen.ordo.spring; + +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; +import org.flywaydb.core.Flyway; +import org.flywaydb.core.api.configuration.FluentConfiguration; +import org.springframework.beans.factory.FactoryBean; +import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.AutoConfigureAfter; +import org.springframework.boot.autoconfigure.AutoConfigureBefore; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.core.env.Environment; + +import javax.sql.DataSource; +import java.lang.reflect.Proxy; + +/** + * Points Flyway at the dialect migration directory before Flyway runs. + * Host {@code spring.flyway.locations} is left unchanged when set. + * + *

Does not inject {@link DataSource} or {@link SqlDialect} at bean-creation time + * (that would cycle with DataSource → Flyway → customizer). Dialect is resolved inside + * {@code customize} from the FluentConfiguration DataSource / {@code ordo.jdbc.dialect}. + * + *

FlywayConfigurationCustomizer moved between Boot 3 and Boot 4; a reflective + * {@link FactoryBean} supplies a proxy for whichever type is on the classpath. + */ +@AutoConfiguration +@ConditionalOnProperty(prefix = "ordo", name = "enabled", havingValue = "true", matchIfMissing = true) +@ConditionalOnClass({DataSource.class, Flyway.class}) +@ConditionalOnBean(DataSource.class) +@AutoConfigureAfter(name = { + "org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration", + "org.springframework.boot.jdbc.autoconfigure.DataSourceAutoConfiguration" +}) +@AutoConfigureBefore(name = { + "org.springframework.boot.autoconfigure.flyway.FlywayAutoConfiguration", + "org.springframework.boot.flyway.autoconfigure.FlywayAutoConfiguration" +}) +@EnableConfigurationProperties(OrdoProperties.class) +public class OrdoFlywayAutoConfiguration { + + private static final String[] FLYWAY_CUSTOMIZER_TYPES = { + "org.springframework.boot.flyway.autoconfigure.FlywayConfigurationCustomizer", + "org.springframework.boot.autoconfigure.flyway.FlywayConfigurationCustomizer" + }; + + @Bean + @ConditionalOnClass(name = "org.flywaydb.core.api.configuration.FluentConfiguration") + public FactoryBean ordoFlywayConfigurationCustomizer( + OrdoProperties ordoProperties, Environment environment) { + Class customizerType = resolveFlywayCustomizerType(); + if (customizerType == null) { + return null; + } + return new FlywayLocationsCustomizerFactoryBean(customizerType, ordoProperties, environment); + } + + static Class resolveFlywayCustomizerType() { + ClassLoader classLoader = OrdoFlywayAutoConfiguration.class.getClassLoader(); + for (String name : FLYWAY_CUSTOMIZER_TYPES) { + try { + return Class.forName(name, false, classLoader); + } catch (ClassNotFoundException ignored) { + // Boot 3 vs Boot 4 + } + } + return null; + } + + private static final class FlywayLocationsCustomizerFactoryBean implements FactoryBean { + private final Class customizerType; + private final OrdoProperties ordoProperties; + private final Environment environment; + + private FlywayLocationsCustomizerFactoryBean(Class customizerType, OrdoProperties ordoProperties, + Environment environment) { + this.customizerType = customizerType; + this.ordoProperties = ordoProperties; + this.environment = environment; + } + + @Override + public Object getObject() { + return Proxy.newProxyInstance(customizerType.getClassLoader(), new Class[] {customizerType}, + (proxy, method, args) -> { + String name = method.getName(); + if ("customize".equals(name) && args != null && args.length == 1) { + applyLocations((FluentConfiguration) args[0]); + return null; + } + if ("equals".equals(name)) { + return proxy == args[0]; + } + if ("hashCode".equals(name)) { + return System.identityHashCode(proxy); + } + if ("toString".equals(name)) { + return "OrdoFlywayLocationsCustomizer"; + } + throw new UnsupportedOperationException(method.toString()); + }); + } + + private void applyLocations(FluentConfiguration configuration) { + if (environment.containsProperty("spring.flyway.locations")) { + return; + } + SqlDialect dialect = SqlDialects.resolve(configuration.getDataSource(), + ordoProperties.getJdbc().getDialect()); + configuration.locations(dialect.flywayLocations()); + } + + @Override + public Class getObjectType() { + return customizerType; + } + } +} diff --git a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java index 18083b6..e418809 100644 --- a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java +++ b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java @@ -19,6 +19,8 @@ import com.jetlumen.ordo.storage.jdbc.JdbcProcessDefinitionRepository; import com.jetlumen.ordo.storage.jdbc.JdbcProcessHistoryRepository; import com.jetlumen.ordo.storage.jdbc.JdbcProcessInstanceRepository; import com.jetlumen.ordo.storage.jdbc.JdbcTransactionExecutor; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.AutoConfigureAfter; @@ -49,7 +51,8 @@ import java.time.Clock; "org.springframework.boot.autoconfigure.flyway.FlywayAutoConfiguration", // Spring Boot 4.x locations (JDBC/Flyway autoconfiguration moved to dedicated modules) "org.springframework.boot.jdbc.autoconfigure.DataSourceAutoConfiguration", - "org.springframework.boot.flyway.autoconfigure.FlywayAutoConfiguration" + "org.springframework.boot.flyway.autoconfigure.FlywayAutoConfiguration", + "com.jetlumen.ordo.spring.OrdoFlywayAutoConfiguration" }) @EnableConfigurationProperties(OrdoProperties.class) public class OrdoJdbcAutoConfiguration { @@ -78,6 +81,12 @@ public class OrdoJdbcAutoConfiguration { return ActionHandler.noop(); } + @Bean + @ConditionalOnMissingBean + public SqlDialect ordoSqlDialect(DataSource dataSource, OrdoProperties ordoProperties) { + return SqlDialects.resolve(dataSource, ordoProperties.getJdbc().getDialect()); + } + @Bean @ConditionalOnMissingBean @DependsOnDatabaseInitialization @@ -93,32 +102,37 @@ public class OrdoJdbcAutoConfiguration { @Bean @ConditionalOnMissingBean - public ProcessDefinitionRepository ordoProcessDefinitionRepository(JdbcConnectionProvider connectionProvider) { - return new JdbcProcessDefinitionRepository(connectionProvider); + public ProcessDefinitionRepository ordoProcessDefinitionRepository( + JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) { + return new JdbcProcessDefinitionRepository(connectionProvider, ordoSqlDialect); } @Bean @ConditionalOnMissingBean - public ProcessInstanceRepository ordoProcessInstanceRepository(JdbcConnectionProvider connectionProvider) { - return new JdbcProcessInstanceRepository(connectionProvider); + public ProcessInstanceRepository ordoProcessInstanceRepository( + JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) { + return new JdbcProcessInstanceRepository(connectionProvider, ordoSqlDialect); } @Bean @ConditionalOnMissingBean - public ApprovalTaskRepository ordoApprovalTaskRepository(JdbcConnectionProvider connectionProvider) { - return new JdbcApprovalTaskRepository(connectionProvider); + public ApprovalTaskRepository ordoApprovalTaskRepository( + JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) { + return new JdbcApprovalTaskRepository(connectionProvider, ordoSqlDialect); } @Bean @ConditionalOnMissingBean - public ProcessHistoryRepository ordoProcessHistoryRepository(JdbcConnectionProvider connectionProvider) { - return new JdbcProcessHistoryRepository(connectionProvider); + public ProcessHistoryRepository ordoProcessHistoryRepository( + JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) { + return new JdbcProcessHistoryRepository(connectionProvider, ordoSqlDialect); } @Bean @ConditionalOnMissingBean - public ActionExecutionRepository ordoActionExecutionRepository(JdbcConnectionProvider connectionProvider) { - return new JdbcActionExecutionRepository(connectionProvider); + public ActionExecutionRepository ordoActionExecutionRepository( + JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) { + return new JdbcActionExecutionRepository(connectionProvider, ordoSqlDialect); } @Bean diff --git a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoProperties.java b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoProperties.java index 1eb9384..8f1e856 100644 --- a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoProperties.java +++ b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoProperties.java @@ -11,6 +11,7 @@ public class OrdoProperties { private final Definitions definitions = new Definitions(); private final Due due = new Due(); + private final Jdbc jdbc = new Jdbc(); public boolean isEnabled() { return enabled; @@ -28,6 +29,10 @@ public class OrdoProperties { return due; } + public Jdbc getJdbc() { + return jdbc; + } + public static class Definitions { private String location = "classpath*:ordo/*.json"; @@ -52,4 +57,17 @@ public class OrdoProperties { this.pollMs = pollMs; } } + + public static class Jdbc { + /** Explicit dialect id ({@code postgresql}, {@code mysql}). Empty means detect from the DataSource. */ + private String dialect; + + public String getDialect() { + return dialect; + } + + public void setDialect(String dialect) { + this.dialect = dialect; + } + } } diff --git a/ordo-spring-boot-autoconfigure/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/ordo-spring-boot-autoconfigure/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index f5adb01..588fcc1 100644 --- a/ordo-spring-boot-autoconfigure/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/ordo-spring-boot-autoconfigure/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -1 +1,2 @@ +com.jetlumen.ordo.spring.OrdoFlywayAutoConfiguration com.jetlumen.ordo.spring.OrdoJdbcAutoConfiguration diff --git a/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java b/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java index 3779947..3a5b05b 100644 --- a/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java +++ b/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java @@ -13,6 +13,7 @@ import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.RoutingCondition; import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.autoconfigure.flyway.FlywayAutoConfiguration; @@ -34,7 +35,8 @@ class OrdoJdbcAutoConfigurationTest { private final ApplicationContextRunner withDataSourceRunner = new ApplicationContextRunner() .withConfiguration(AutoConfigurations.of( - DataSourceAutoConfiguration.class, FlywayAutoConfiguration.class, OrdoJdbcAutoConfiguration.class)) + DataSourceAutoConfiguration.class, OrdoFlywayAutoConfiguration.class, + FlywayAutoConfiguration.class, OrdoJdbcAutoConfiguration.class)) .withPropertyValues( "spring.datasource.url=jdbc:h2:mem:ordo_" + UUID.randomUUID() + ";MODE=PostgreSQL;DATABASE_TO_LOWER=TRUE;DB_CLOSE_DELAY=-1", @@ -47,6 +49,7 @@ class OrdoJdbcAutoConfigurationTest { void assemblesJdbcBackedEngineWhenDataSourceIsPresent() { withDataSourceRunner.run(context -> { assertThat(context).hasSingleBean(OrdoEngine.class); + assertThat(context.getBean(SqlDialect.class).id()).isEqualTo("postgresql"); OrdoEngine engine = context.getBean(OrdoEngine.class); engine.publish(LEAVE_REQUEST); @@ -60,6 +63,12 @@ class OrdoJdbcAutoConfigurationTest { }); } + @Test + void honoursExplicitMysqlDialect() { + withDataSourceRunner.withPropertyValues("ordo.jdbc.dialect=mysql", "spring.flyway.enabled=false") + .run(context -> assertThat(context.getBean(SqlDialect.class).id()).isEqualTo("mysql")); + } + @Test void backsOffWhenNoDataSourceBeanIsPresent() { withoutDataSourceRunner.run(context -> diff --git a/ordo-spring-boot-starter/pom.xml b/ordo-spring-boot-starter/pom.xml index 8cb0735..6063a19 100644 --- a/ordo-spring-boot-starter/pom.xml +++ b/ordo-spring-boot-starter/pom.xml @@ -47,16 +47,22 @@ org.flywaydb flyway-database-postgresql + + org.flywaydb + flyway-mysql + diff --git a/ordo-storage-jdbc/pom.xml b/ordo-storage-jdbc/pom.xml index 245158b..b4322b8 100644 --- a/ordo-storage-jdbc/pom.xml +++ b/ordo-storage-jdbc/pom.xml @@ -32,6 +32,11 @@ postgresql test + + com.mysql + mysql-connector-j + test + org.junit.jupiter junit-jupiter diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java index 0e95127..4f930cf 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java @@ -5,6 +5,8 @@ import com.jetlumen.ordo.api.ActionExecutionStatus; import com.jetlumen.ordo.api.query.Page; import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ActionExecutionRepository; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import com.jetlumen.ordo.storage.jdbc.mapper.ActionExecutionMapper; import java.sql.Connection; @@ -28,9 +30,15 @@ public final class JdbcActionExecutionRepository implements ActionExecutionRepos + " WHERE id = ? AND status = 'PENDING'"; private final JdbcConnectionProvider connectionProvider; + private final SqlDialect dialect; public JdbcActionExecutionRepository(JdbcConnectionProvider connectionProvider) { + this(connectionProvider, SqlDialects.postgresql()); + } + + public JdbcActionExecutionRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); } @Override @@ -71,8 +79,8 @@ public final class JdbcActionExecutionRepository implements ActionExecutionRepos try { long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_action_execution WHERE instance_id = ?", List.of(instanceId)); - String sql = "SELECT " + COLUMNS + " FROM ordo_action_execution WHERE instance_id = ?" - + " ORDER BY started_at ASC, id ASC LIMIT ? OFFSET ?"; + String sql = dialect.limit("SELECT " + COLUMNS + " FROM ordo_action_execution WHERE instance_id = ?" + + " ORDER BY started_at ASC, id ASC"); List content = new ArrayList<>(); try (PreparedStatement select = connection.prepareStatement(sql)) { select.setString(1, instanceId); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java index eeb7b61..fad4529 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java @@ -6,6 +6,8 @@ import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.query.TaskQuery; import com.jetlumen.ordo.api.repository.ApprovalTaskRepository; import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper; import java.sql.Connection; @@ -47,17 +49,25 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository private static final String CLAIM_IF_DUE = "UPDATE ordo_approval_task SET due_at = NULL WHERE id = ? AND status = 'PENDING' AND assignee = ?" + " AND due_at IS NOT NULL AND due_at <= ?"; - private static final String SELECT_DUE_PENDING = + private static final String SELECT_DUE_PENDING_BASE = "SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL" - + " AND due_at <= ? ORDER BY due_at, id LIMIT ?"; + + " AND due_at <= ? ORDER BY due_at, id"; private static final String TASK_COLUMNS_QUALIFIED = "t.id, t.instance_id, t.step_id, t.task_name, t.assignee, t.status, t.created_at, t.completed_at," + " t.action_actor, t.action_comment, t.action_at, t.due_at"; private final JdbcConnectionProvider connectionProvider; + private final SqlDialect dialect; + private final String selectDuePending; public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider) { + this(connectionProvider, SqlDialects.postgresql()); + } + + public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); + this.selectDuePending = dialect.limit(SELECT_DUE_PENDING_BASE, false); } @Override @@ -153,7 +163,7 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository throw new IllegalArgumentException("limit must be positive"); } Connection connection = connectionProvider.getConnection(); - try (PreparedStatement select = connection.prepareStatement(SELECT_DUE_PENDING)) { + try (PreparedStatement select = connection.prepareStatement(selectDuePending)) { select.setTimestamp(1, Timestamp.from(now)); select.setInt(2, limit); try (ResultSet resultSet = select.executeQuery()) { @@ -197,8 +207,8 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository try { long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_approval_task t" + where.sql(), where.params()); - String sql = "SELECT " + TASK_COLUMNS_QUALIFIED + " FROM ordo_approval_task t" + where.sql() - + " ORDER BY t.created_at DESC, t.id DESC LIMIT ? OFFSET ?"; + String sql = dialect.limit("SELECT " + TASK_COLUMNS_QUALIFIED + " FROM ordo_approval_task t" + where.sql() + + " ORDER BY t.created_at DESC, t.id DESC"); List content = new ArrayList<>(); try (PreparedStatement select = connection.prepareStatement(sql)) { int index = PageSupport.bindParams(select, where.params()); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java index 303c91f..0aa458f 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java @@ -6,6 +6,8 @@ import com.jetlumen.ordo.api.StepTransition; import com.jetlumen.ordo.api.query.Page; import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper.StepRow; import com.jetlumen.ordo.storage.jdbc.mapper.StepTransitionMapper; @@ -30,8 +32,8 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR "INSERT INTO ordo_process (id, current_version, name) VALUES (?, ?, ?)"; private static final String UPDATE_PROCESS = "UPDATE ordo_process SET current_version = ?, name = ? WHERE id = ?"; - private static final String LOCK_PROCESS = - "SELECT current_version FROM ordo_process WHERE id = ? FOR UPDATE"; + private static final String LOCK_PROCESS_BASE = + "SELECT current_version FROM ordo_process WHERE id = ?"; private static final String INSERT_DEFINITION = "INSERT INTO ordo_process_definition (id, version, name, created_at) VALUES (?, ?, ?, ?)"; private static final String INSERT_STEP = @@ -56,20 +58,32 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR private static final String SELECT_TRANSITIONS = "SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition " + "WHERE definition_id = ? AND definition_version = ? ORDER BY from_step_id, priority"; - private static final String SELECT_LATEST_PAGE = + private static final String SELECT_LATEST_PAGE_BASE = "SELECT p.id, p.current_version, d.name FROM ordo_process p" + " JOIN ordo_process_definition d ON d.id = p.id AND d.version = p.current_version" - + " ORDER BY p.id LIMIT ? OFFSET ?"; + + " ORDER BY p.id"; private static final String COUNT_PROCESSES = "SELECT COUNT(*) FROM ordo_process"; - private static final String SELECT_VERSIONS_PAGE = - "SELECT version, name FROM ordo_process_definition WHERE id = ? ORDER BY version DESC LIMIT ? OFFSET ?"; + private static final String SELECT_VERSIONS_PAGE_BASE = + "SELECT version, name FROM ordo_process_definition WHERE id = ? ORDER BY version DESC"; private static final String COUNT_VERSIONS = "SELECT COUNT(*) FROM ordo_process_definition WHERE id = ?"; private final JdbcConnectionProvider connectionProvider; + private final SqlDialect dialect; + private final String lockProcess; + private final String selectLatestPage; + private final String selectVersionsPage; public JdbcProcessDefinitionRepository(JdbcConnectionProvider connectionProvider) { + this(connectionProvider, SqlDialects.postgresql()); + } + + public JdbcProcessDefinitionRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); + this.lockProcess = dialect.forUpdate(LOCK_PROCESS_BASE); + this.selectLatestPage = dialect.limit(SELECT_LATEST_PAGE_BASE); + this.selectVersionsPage = dialect.limit(SELECT_VERSIONS_PAGE_BASE); } @Override @@ -147,7 +161,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR try { long total = PageSupport.count(connection, COUNT_PROCESSES, List.of()); List content = new ArrayList<>(); - try (PreparedStatement select = connection.prepareStatement(SELECT_LATEST_PAGE)) { + try (PreparedStatement select = connection.prepareStatement(selectLatestPage)) { select.setInt(1, pageRequest.size()); select.setInt(2, pageRequest.offset()); List versions = new ArrayList<>(); @@ -177,7 +191,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR try { long total = PageSupport.count(connection, COUNT_VERSIONS, List.of(definitionId)); List versionNumbers = new ArrayList<>(); - try (PreparedStatement select = connection.prepareStatement(SELECT_VERSIONS_PAGE)) { + try (PreparedStatement select = connection.prepareStatement(selectVersionsPage)) { select.setString(1, definitionId); select.setInt(2, pageRequest.size()); select.setInt(3, pageRequest.offset()); @@ -254,8 +268,8 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR return new ProcessDefinition(definitionId, version, name, steps, transitions); } - private static Integer lockCurrentVersion(Connection connection, String definitionId) throws SQLException { - try (PreparedStatement select = connection.prepareStatement(LOCK_PROCESS)) { + private Integer lockCurrentVersion(Connection connection, String definitionId) throws SQLException { + try (PreparedStatement select = connection.prepareStatement(lockProcess)) { select.setString(1, definitionId); try (ResultSet resultSet = select.executeQuery()) { return resultSet.next() ? resultSet.getInt("current_version") : null; @@ -284,7 +298,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR } } - private static boolean insertProcessRow(Connection connection, String id, int version, String name) + private boolean insertProcessRow(Connection connection, String id, int version, String name) throws SQLException { try (PreparedStatement insert = connection.prepareStatement(INSERT_PROCESS)) { insert.setString(1, id); @@ -293,7 +307,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR insert.executeUpdate(); return true; } catch (SQLException e) { - if (isDuplicateKey(e)) { + if (dialect.isDuplicateKey(e)) { return false; } throw e; @@ -370,7 +384,4 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR } } - private static boolean isDuplicateKey(SQLException e) { - return "23505".equals(e.getSQLState()); - } } diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java index b135eee..fc3b305 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java @@ -4,6 +4,8 @@ import com.jetlumen.ordo.api.ProcessEvent; import com.jetlumen.ordo.api.query.Page; import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessHistoryRepository; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import com.jetlumen.ordo.storage.jdbc.mapper.ProcessEventMapper; import java.sql.Connection; @@ -23,9 +25,15 @@ public final class JdbcProcessHistoryRepository implements ProcessHistoryReposit + " VALUES (?, ?, ?, ?, ?, ?, ?, ?)"; private final JdbcConnectionProvider connectionProvider; + private final SqlDialect dialect; public JdbcProcessHistoryRepository(JdbcConnectionProvider connectionProvider) { + this(connectionProvider, SqlDialects.postgresql()); + } + + public JdbcProcessHistoryRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); } @Override @@ -50,8 +58,8 @@ public final class JdbcProcessHistoryRepository implements ProcessHistoryReposit try { long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_process_event WHERE instance_id = ?", List.of(instanceId)); - String sql = "SELECT " + COLUMNS + " FROM ordo_process_event WHERE instance_id = ?" - + " ORDER BY occurred_at ASC, id ASC LIMIT ? OFFSET ?"; + String sql = dialect.limit("SELECT " + COLUMNS + " FROM ordo_process_event WHERE instance_id = ?" + + " ORDER BY occurred_at ASC, id ASC"); List content = new ArrayList<>(); try (PreparedStatement select = connection.prepareStatement(sql)) { select.setString(1, instanceId); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java index 5458691..356f18f 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java @@ -7,6 +7,8 @@ import com.jetlumen.ordo.api.query.Page; import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper; import java.sql.Connection; @@ -31,13 +33,21 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos private static final String SELECT_INSTANCE = "SELECT id, definition_id, definition_version, initiator, status, context_json, started_at, finished_at" + " FROM ordo_process_instance WHERE id = ?"; - private static final String EXISTS_RUNNING = - "SELECT 1 FROM ordo_process_instance WHERE definition_id = ? AND status = ? LIMIT 1"; + private static final String EXISTS_RUNNING_BASE = + "SELECT 1 FROM ordo_process_instance WHERE definition_id = ? AND status = ?"; private final JdbcConnectionProvider connectionProvider; + private final SqlDialect dialect; + private final String existsRunning; public JdbcProcessInstanceRepository(JdbcConnectionProvider connectionProvider) { + this(connectionProvider, SqlDialects.postgresql()); + } + + public JdbcProcessInstanceRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); + this.existsRunning = dialect.limit(EXISTS_RUNNING_BASE, false); } @Override @@ -88,9 +98,10 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos @Override public boolean existsRunning(String definitionId) { Connection connection = connectionProvider.getConnection(); - try (PreparedStatement select = connection.prepareStatement(EXISTS_RUNNING)) { + try (PreparedStatement select = connection.prepareStatement(existsRunning)) { select.setString(1, definitionId); select.setString(2, ProcessStatus.RUNNING.name()); + select.setInt(3, 1); try (ResultSet resultSet = select.executeQuery()) { return resultSet.next(); } @@ -124,9 +135,9 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos try { long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_process_instance" + where.sql(), where.params()); - String sql = "SELECT id, definition_id, definition_version, initiator, status, context_json, started_at," - + " finished_at FROM ordo_process_instance" + where.sql() - + " ORDER BY started_at DESC, id DESC LIMIT ? OFFSET ?"; + String sql = dialect.limit("SELECT id, definition_id, definition_version, initiator, status, context_json," + + " started_at, finished_at FROM ordo_process_instance" + where.sql() + + " ORDER BY started_at DESC, id DESC"); List content = new ArrayList<>(); try (PreparedStatement select = connection.prepareStatement(sql)) { int index = PageSupport.bindParams(select, where.params()); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/LimitOffsetSqlDialect.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/LimitOffsetSqlDialect.java new file mode 100644 index 0000000..5f11a23 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/LimitOffsetSqlDialect.java @@ -0,0 +1,19 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +abstract class LimitOffsetSqlDialect implements SqlDialect { + + @Override + public final String limit(String sql) { + return limit(sql, true); + } + + @Override + public final String limit(String sql, boolean offset) { + return offset ? sql + " LIMIT ? OFFSET ?" : sql + " LIMIT ?"; + } + + @Override + public final String forUpdate(String sql) { + return sql + " FOR UPDATE"; + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/MysqlSqlDialect.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/MysqlSqlDialect.java new file mode 100644 index 0000000..d3f688e --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/MysqlSqlDialect.java @@ -0,0 +1,38 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +import java.sql.DatabaseMetaData; +import java.sql.SQLException; +import java.util.Locale; + +/** MySQL / MariaDB (and H2 {@code MODE=MySQL}) dialect. */ +public final class MysqlSqlDialect extends LimitOffsetSqlDialect { + + static final String ID = "mysql"; + + @Override + public String id() { + return ID; + } + + @Override + public boolean supports(DatabaseMetaData metaData) throws SQLException { + String product = metaData.getDatabaseProductName(); + if (product != null) { + String lower = product.toLowerCase(Locale.ROOT); + if (lower.contains("mysql") || lower.contains("mariadb")) { + return true; + } + } + return PostgresSqlDialect.isH2Mode(metaData, "MYSQL"); + } + + @Override + public boolean isDuplicateKey(SQLException exception) { + return "23000".equals(exception.getSQLState()) && exception.getErrorCode() == 1062; + } + + @Override + public String[] flywayLocations() { + return new String[] {"classpath:db/mysql/migration"}; + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/PostgresSqlDialect.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/PostgresSqlDialect.java new file mode 100644 index 0000000..25d4fff --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/PostgresSqlDialect.java @@ -0,0 +1,52 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +import java.sql.DatabaseMetaData; +import java.sql.SQLException; +import java.util.Locale; + +/** PostgreSQL (and H2 {@code MODE=PostgreSQL}) dialect. */ +public final class PostgresSqlDialect extends LimitOffsetSqlDialect { + + static final String ID = "postgresql"; + + @Override + public String id() { + return ID; + } + + @Override + public boolean supports(DatabaseMetaData metaData) throws SQLException { + String product = metaData.getDatabaseProductName(); + if (product != null && product.toLowerCase(Locale.ROOT).contains("postgresql")) { + return true; + } + return isH2Mode(metaData, "POSTGRESQL"); + } + + @Override + public boolean isDuplicateKey(SQLException exception) { + return "23505".equals(exception.getSQLState()); + } + + @Override + public String[] flywayLocations() { + return new String[] {"classpath:db/postgresql/migration"}; + } + + static boolean isH2Mode(DatabaseMetaData metaData, String mode) throws SQLException { + String product = metaData.getDatabaseProductName(); + if (product == null || !product.equalsIgnoreCase("H2")) { + return false; + } + String url = metaData.getURL(); + if (url != null && url.toUpperCase(Locale.ROOT).contains("MODE=" + mode)) { + return true; + } + // H2's DatabaseMetaData.getURL() often omits MODE=; read the runtime setting instead. + try (var statement = metaData.getConnection().createStatement(); + var resultSet = statement.executeQuery( + "SELECT SETTING_VALUE FROM INFORMATION_SCHEMA.SETTINGS WHERE SETTING_NAME = 'MODE'")) { + return resultSet.next() && mode.equalsIgnoreCase(resultSet.getString(1)); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialect.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialect.java new file mode 100644 index 0000000..ad9c92c --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialect.java @@ -0,0 +1,29 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +import java.sql.DatabaseMetaData; +import java.sql.SQLException; + +/** + * Database-specific SQL fragments and Flyway locations for JDBC storage. + * Additional dialects can be supplied via {@link java.util.ServiceLoader}. + */ +public interface SqlDialect { + + String id(); + + boolean supports(DatabaseMetaData metaData) throws SQLException; + + boolean isDuplicateKey(SQLException exception); + + /** Appends {@code LIMIT ? OFFSET ?}. */ + String limit(String sql); + + /** + * Appends a row limit. When {@code offset} is {@code true}, also binds {@code OFFSET ?}. + */ + String limit(String sql, boolean offset); + + String forUpdate(String sql); + + String[] flywayLocations(); +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialects.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialects.java new file mode 100644 index 0000000..bb5e95a --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialects.java @@ -0,0 +1,78 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +import javax.sql.DataSource; +import java.sql.Connection; +import java.sql.DatabaseMetaData; +import java.sql.SQLException; +import java.util.ArrayList; +import java.util.List; +import java.util.Locale; +import java.util.Objects; +import java.util.ServiceLoader; +import java.util.stream.Collectors; + +/** Loads {@link SqlDialect} implementations and resolves one for a DataSource or explicit id. */ +public final class SqlDialects { + + private SqlDialects() { + } + + public static List load() { + List dialects = new ArrayList<>(); + ServiceLoader.load(SqlDialect.class, SqlDialect.class.getClassLoader()).forEach(dialects::add); + if (dialects.isEmpty()) { + dialects.add(new PostgresSqlDialect()); + dialects.add(new MysqlSqlDialect()); + } + return List.copyOf(dialects); + } + + public static SqlDialect postgresql() { + return new PostgresSqlDialect(); + } + + public static SqlDialect mysql() { + return new MysqlSqlDialect(); + } + + /** + * Resolves a dialect. A non-blank {@code dialectId} wins; otherwise the DataSource metadata + * is matched against registered dialects. + */ + public static SqlDialect resolve(DataSource dataSource, String dialectId) { + List dialects = load(); + if (dialectId != null && !dialectId.isBlank()) { + return byId(dialects, dialectId); + } + Objects.requireNonNull(dataSource, "dataSource must not be null"); + try (Connection connection = dataSource.getConnection()) { + return detect(dialects, connection.getMetaData()); + } catch (SQLException e) { + throw new IllegalStateException("failed to detect ordo JDBC dialect", e); + } + } + + static SqlDialect byId(List dialects, String dialectId) { + String wanted = dialectId.trim().toLowerCase(Locale.ROOT); + return dialects.stream() + .filter(dialect -> dialect.id().equalsIgnoreCase(wanted)) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException( + "unknown ordo JDBC dialect '" + dialectId + "'; known: " + ids(dialects))); + } + + static SqlDialect detect(List dialects, DatabaseMetaData metaData) throws SQLException { + for (SqlDialect dialect : dialects) { + if (dialect.supports(metaData)) { + return dialect; + } + } + String product = metaData.getDatabaseProductName(); + throw new IllegalStateException( + "no ordo JDBC dialect for " + product + "; known: " + ids(dialects)); + } + + private static String ids(List dialects) { + return dialects.stream().map(SqlDialect::id).collect(Collectors.joining(", ")); + } +} diff --git a/ordo-storage-jdbc/src/main/resources/META-INF/services/com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect b/ordo-storage-jdbc/src/main/resources/META-INF/services/com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect new file mode 100644 index 0000000..c1d475b --- /dev/null +++ b/ordo-storage-jdbc/src/main/resources/META-INF/services/com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect @@ -0,0 +1,2 @@ +com.jetlumen.ordo.storage.jdbc.dialect.PostgresSqlDialect +com.jetlumen.ordo.storage.jdbc.dialect.MysqlSqlDialect diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql deleted file mode 100644 index bd9d011..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql +++ /dev/null @@ -1,48 +0,0 @@ --- Ordo approval workflow tables (v1). Target database: PostgreSQL. --- The DDL sticks to portable types so the same script also runs on H2 in --- PostgreSQL compatibility mode, which the integration tests use. - -CREATE TABLE ordo_process_definition ( - id VARCHAR(64) PRIMARY KEY, - name VARCHAR(255) NOT NULL -); - -CREATE TABLE ordo_approval_step ( - definition_id VARCHAR(64) NOT NULL, - step_id VARCHAR(64) NOT NULL, - step_name VARCHAR(255) NOT NULL, - assignee VARCHAR(255) NOT NULL, - step_order INTEGER NOT NULL, - PRIMARY KEY (definition_id, step_id), - CONSTRAINT fk_approval_step_definition FOREIGN KEY (definition_id) REFERENCES ordo_process_definition (id) -); - -CREATE TABLE ordo_process_instance ( - id VARCHAR(36) PRIMARY KEY, - definition_id VARCHAR(64) NOT NULL, - initiator VARCHAR(255) NOT NULL, - status VARCHAR(32) NOT NULL, - context_json TEXT NOT NULL, - started_at TIMESTAMP NOT NULL, - finished_at TIMESTAMP, - CONSTRAINT fk_process_instance_definition FOREIGN KEY (definition_id) REFERENCES ordo_process_definition (id) -); - -CREATE TABLE ordo_approval_task ( - id VARCHAR(36) PRIMARY KEY, - instance_id VARCHAR(36) NOT NULL, - step_id VARCHAR(64) NOT NULL, - task_name VARCHAR(255) NOT NULL, - assignee VARCHAR(255) NOT NULL, - status VARCHAR(32) NOT NULL, - created_at TIMESTAMP NOT NULL, - completed_at TIMESTAMP, - action_actor VARCHAR(255), - action_comment TEXT, - action_at TIMESTAMP, - CONSTRAINT fk_approval_task_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) -); - -CREATE INDEX idx_approval_task_instance ON ordo_approval_task (instance_id); -CREATE INDEX idx_approval_task_status_assignee ON ordo_approval_task (status, assignee); -CREATE INDEX idx_process_instance_definition ON ordo_process_instance (definition_id); diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V2__add_step_candidates_and_policy.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V2__add_step_candidates_and_policy.sql deleted file mode 100644 index f1c7cdf..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V2__add_step_candidates_and_policy.sql +++ /dev/null @@ -1,23 +0,0 @@ --- Adds multi-candidate (any/all) approval step support. --- A step no longer has a single fixed assignee; instead it has one or more --- candidates in ordo_step_candidate, and a policy column decides whether any --- single candidate approval is enough (ANY) or every candidate must approve (ALL). - -ALTER TABLE ordo_approval_step ADD COLUMN policy VARCHAR(16) NOT NULL DEFAULT 'ANY'; - -CREATE TABLE ordo_step_candidate ( - definition_id VARCHAR(64) NOT NULL, - step_id VARCHAR(64) NOT NULL, - candidate VARCHAR(255) NOT NULL, - candidate_order INTEGER NOT NULL, - PRIMARY KEY (definition_id, step_id, candidate), - CONSTRAINT fk_step_candidate_step FOREIGN KEY (definition_id, step_id) - REFERENCES ordo_approval_step (definition_id, step_id) -); - --- Migrate any existing single-assignee steps into the new candidate table before --- the now-unused column is dropped. -INSERT INTO ordo_step_candidate (definition_id, step_id, candidate, candidate_order) -SELECT definition_id, step_id, assignee, 0 FROM ordo_approval_step; - -ALTER TABLE ordo_approval_step DROP COLUMN assignee; diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V3__add_step_transitions.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V3__add_step_transitions.sql deleted file mode 100644 index 02742be..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V3__add_step_transitions.sql +++ /dev/null @@ -1,12 +0,0 @@ -CREATE TABLE ordo_step_transition ( - definition_id VARCHAR(64) NOT NULL, - from_step_id VARCHAR(64) NOT NULL, - to_step_id VARCHAR(64), - condition_key VARCHAR(255), - priority INTEGER NOT NULL, - PRIMARY KEY (definition_id, from_step_id, priority), - CONSTRAINT fk_transition_from FOREIGN KEY (definition_id, from_step_id) - REFERENCES ordo_approval_step (definition_id, step_id), - CONSTRAINT fk_transition_to FOREIGN KEY (definition_id, to_step_id) - REFERENCES ordo_approval_step (definition_id, step_id) -); diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V4__add_step_kind_and_action_key.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V4__add_step_kind_and_action_key.sql deleted file mode 100644 index 23e628b..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V4__add_step_kind_and_action_key.sql +++ /dev/null @@ -1,3 +0,0 @@ --- Adds ACTION step kind and an optional action_key (host ActionHandler lookup). -ALTER TABLE ordo_approval_step ADD COLUMN kind VARCHAR(16) NOT NULL DEFAULT 'APPROVAL'; -ALTER TABLE ordo_approval_step ADD COLUMN action_key VARCHAR(255); diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql deleted file mode 100644 index 52ed75d..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql +++ /dev/null @@ -1,29 +0,0 @@ --- Process history events and ACTION-step execution records. - -CREATE TABLE ordo_process_event ( - id VARCHAR(36) PRIMARY KEY, - instance_id VARCHAR(36) NOT NULL, - task_id VARCHAR(36), - step_id VARCHAR(64), - event_type VARCHAR(32) NOT NULL, - actor VARCHAR(255), - detail TEXT, - occurred_at TIMESTAMP NOT NULL, - CONSTRAINT fk_process_event_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) -); - -CREATE INDEX idx_process_event_instance_time ON ordo_process_event (instance_id, occurred_at, id); - -CREATE TABLE ordo_action_execution ( - id VARCHAR(36) PRIMARY KEY, - instance_id VARCHAR(36) NOT NULL, - step_id VARCHAR(64) NOT NULL, - action_key VARCHAR(255) NOT NULL, - status VARCHAR(32) NOT NULL, - error_message TEXT, - started_at TIMESTAMP NOT NULL, - finished_at TIMESTAMP, - CONSTRAINT fk_action_execution_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) -); - -CREATE INDEX idx_action_execution_instance ON ordo_action_execution (instance_id); diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V6__add_step_due.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V6__add_step_due.sql deleted file mode 100644 index d79a898..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V6__add_step_due.sql +++ /dev/null @@ -1,8 +0,0 @@ --- SLA / due escalation: task due_at plus step due policy columns. -ALTER TABLE ordo_approval_task ADD COLUMN due_at TIMESTAMP; -CREATE INDEX idx_approval_task_status_due ON ordo_approval_task (status, due_at); - -ALTER TABLE ordo_approval_step ADD COLUMN due_after VARCHAR(32); -ALTER TABLE ordo_approval_step ADD COLUMN due_then VARCHAR(16); -ALTER TABLE ordo_approval_step ADD COLUMN due_to VARCHAR(255); -ALTER TABLE ordo_approval_step ADD COLUMN due_action VARCHAR(255); diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V7__definition_versions.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V7__definition_versions.sql deleted file mode 100644 index d0af913..0000000 --- a/ordo-storage-jdbc/src/main/resources/db/migration/V7__definition_versions.sql +++ /dev/null @@ -1,96 +0,0 @@ -CREATE TABLE ordo_process ( - id VARCHAR(64) PRIMARY KEY, - current_version INTEGER NOT NULL, - name VARCHAR(255) NOT NULL -); - -INSERT INTO ordo_process (id, current_version, name) -SELECT id, 1, name FROM ordo_process_definition; - -ALTER TABLE ordo_process_instance DROP CONSTRAINT fk_process_instance_definition; -ALTER TABLE ordo_approval_step DROP CONSTRAINT fk_approval_step_definition; -ALTER TABLE ordo_step_candidate DROP CONSTRAINT fk_step_candidate_step; -ALTER TABLE ordo_step_transition DROP CONSTRAINT fk_transition_from; -ALTER TABLE ordo_step_transition DROP CONSTRAINT fk_transition_to; - -CREATE TABLE ordo_process_definition_v7 ( - id VARCHAR(64) NOT NULL, - version INTEGER NOT NULL, - name VARCHAR(255) NOT NULL, - created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, - PRIMARY KEY (id, version) -); - -INSERT INTO ordo_process_definition_v7 (id, version, name) -SELECT id, 1, name FROM ordo_process_definition; - -DROP TABLE ordo_process_definition; -ALTER TABLE ordo_process_definition_v7 RENAME TO ordo_process_definition; - -CREATE TABLE ordo_approval_step_v7 ( - definition_id VARCHAR(64) NOT NULL, - definition_version INTEGER NOT NULL, - step_id VARCHAR(64) NOT NULL, - step_name VARCHAR(255) NOT NULL, - step_order INTEGER NOT NULL, - policy VARCHAR(16) NOT NULL, - kind VARCHAR(16) NOT NULL, - action_key VARCHAR(255), - due_after VARCHAR(32), - due_then VARCHAR(16), - due_to VARCHAR(255), - due_action VARCHAR(255), - PRIMARY KEY (definition_id, definition_version, step_id), - CONSTRAINT fk_approval_step_definition FOREIGN KEY (definition_id, definition_version) - REFERENCES ordo_process_definition (id, version) -); - -INSERT INTO ordo_approval_step_v7 (definition_id, definition_version, step_id, step_name, step_order, policy, kind, - action_key, due_after, due_then, due_to, due_action) -SELECT definition_id, 1, step_id, step_name, step_order, policy, kind, action_key, due_after, due_then, due_to, - due_action -FROM ordo_approval_step; - -DROP TABLE ordo_approval_step; -ALTER TABLE ordo_approval_step_v7 RENAME TO ordo_approval_step; - -CREATE TABLE ordo_step_candidate_v7 ( - definition_id VARCHAR(64) NOT NULL, - definition_version INTEGER NOT NULL, - step_id VARCHAR(64) NOT NULL, - candidate VARCHAR(255) NOT NULL, - candidate_order INTEGER NOT NULL, - PRIMARY KEY (definition_id, definition_version, step_id, candidate), - CONSTRAINT fk_step_candidate_step FOREIGN KEY (definition_id, definition_version, step_id) - REFERENCES ordo_approval_step (definition_id, definition_version, step_id) -); - -INSERT INTO ordo_step_candidate_v7 (definition_id, definition_version, step_id, candidate, candidate_order) -SELECT definition_id, 1, step_id, candidate, candidate_order FROM ordo_step_candidate; - -DROP TABLE ordo_step_candidate; -ALTER TABLE ordo_step_candidate_v7 RENAME TO ordo_step_candidate; - -CREATE TABLE ordo_step_transition_v7 ( - definition_id VARCHAR(64) NOT NULL, - definition_version INTEGER NOT NULL, - from_step_id VARCHAR(64) NOT NULL, - to_step_id VARCHAR(64), - condition_key VARCHAR(255), - priority INTEGER NOT NULL, - PRIMARY KEY (definition_id, definition_version, from_step_id, priority), - CONSTRAINT fk_transition_from FOREIGN KEY (definition_id, definition_version, from_step_id) - REFERENCES ordo_approval_step (definition_id, definition_version, step_id), - CONSTRAINT fk_transition_to FOREIGN KEY (definition_id, definition_version, to_step_id) - REFERENCES ordo_approval_step (definition_id, definition_version, step_id) -); - -INSERT INTO ordo_step_transition_v7 (definition_id, definition_version, from_step_id, to_step_id, condition_key, priority) -SELECT definition_id, 1, from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition; - -DROP TABLE ordo_step_transition; -ALTER TABLE ordo_step_transition_v7 RENAME TO ordo_step_transition; - -ALTER TABLE ordo_process_instance ADD COLUMN definition_version INTEGER DEFAULT 1 NOT NULL; -ALTER TABLE ordo_process_instance ADD CONSTRAINT fk_process_instance_definition - FOREIGN KEY (definition_id, definition_version) REFERENCES ordo_process_definition (id, version); diff --git a/ordo-storage-jdbc/src/main/resources/db/mysql/migration/V1__baseline.sql b/ordo-storage-jdbc/src/main/resources/db/mysql/migration/V1__baseline.sql new file mode 100644 index 0000000..d598fd7 --- /dev/null +++ b/ordo-storage-jdbc/src/main/resources/db/mysql/migration/V1__baseline.sql @@ -0,0 +1,125 @@ +-- Ordo schema baseline (MySQL / MariaDB). + +CREATE TABLE ordo_process ( + id VARCHAR(64) NOT NULL, + current_version INTEGER NOT NULL, + name VARCHAR(255) NOT NULL, + PRIMARY KEY (id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_process_definition ( + id VARCHAR(64) NOT NULL, + version INTEGER NOT NULL, + name VARCHAR(255) NOT NULL, + created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), + PRIMARY KEY (id, version) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_approval_step ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + step_id VARCHAR(64) NOT NULL, + step_name VARCHAR(255) NOT NULL, + step_order INTEGER NOT NULL, + policy VARCHAR(16) NOT NULL, + kind VARCHAR(16) NOT NULL, + action_key VARCHAR(255), + due_after VARCHAR(32), + due_then VARCHAR(16), + due_to VARCHAR(255), + due_action VARCHAR(255), + PRIMARY KEY (definition_id, definition_version, step_id), + CONSTRAINT fk_approval_step_definition FOREIGN KEY (definition_id, definition_version) + REFERENCES ordo_process_definition (id, version) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_step_candidate ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + step_id VARCHAR(64) NOT NULL, + candidate VARCHAR(255) NOT NULL, + candidate_order INTEGER NOT NULL, + PRIMARY KEY (definition_id, definition_version, step_id, candidate), + CONSTRAINT fk_step_candidate_step FOREIGN KEY (definition_id, definition_version, step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_step_transition ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + from_step_id VARCHAR(64) NOT NULL, + to_step_id VARCHAR(64), + condition_key VARCHAR(255), + priority INTEGER NOT NULL, + PRIMARY KEY (definition_id, definition_version, from_step_id, priority), + CONSTRAINT fk_transition_from FOREIGN KEY (definition_id, definition_version, from_step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id), + CONSTRAINT fk_transition_to FOREIGN KEY (definition_id, definition_version, to_step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_process_instance ( + id VARCHAR(36) NOT NULL, + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + initiator VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + context_json TEXT NOT NULL, + started_at DATETIME(3) NOT NULL, + finished_at DATETIME(3), + PRIMARY KEY (id), + CONSTRAINT fk_process_instance_definition FOREIGN KEY (definition_id, definition_version) + REFERENCES ordo_process_definition (id, version) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE ordo_approval_task ( + id VARCHAR(36) NOT NULL, + instance_id VARCHAR(36) NOT NULL, + step_id VARCHAR(64) NOT NULL, + task_name VARCHAR(255) NOT NULL, + assignee VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + created_at DATETIME(3) NOT NULL, + completed_at DATETIME(3), + action_actor VARCHAR(255), + action_comment TEXT, + action_at DATETIME(3), + due_at DATETIME(3), + PRIMARY KEY (id), + CONSTRAINT fk_approval_task_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE INDEX idx_approval_task_instance ON ordo_approval_task (instance_id); +CREATE INDEX idx_approval_task_status_assignee ON ordo_approval_task (status, assignee); +CREATE INDEX idx_approval_task_status_due ON ordo_approval_task (status, due_at); +CREATE INDEX idx_process_instance_definition ON ordo_process_instance (definition_id); + +CREATE TABLE ordo_process_event ( + id VARCHAR(36) NOT NULL, + instance_id VARCHAR(36) NOT NULL, + task_id VARCHAR(36), + step_id VARCHAR(64), + event_type VARCHAR(32) NOT NULL, + actor VARCHAR(255), + detail TEXT, + occurred_at DATETIME(3) NOT NULL, + PRIMARY KEY (id), + CONSTRAINT fk_process_event_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE INDEX idx_process_event_instance_time ON ordo_process_event (instance_id, occurred_at, id); + +CREATE TABLE ordo_action_execution ( + id VARCHAR(36) NOT NULL, + instance_id VARCHAR(36) NOT NULL, + step_id VARCHAR(64) NOT NULL, + action_key VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + error_message TEXT, + started_at DATETIME(3) NOT NULL, + finished_at DATETIME(3), + PRIMARY KEY (id), + CONSTRAINT fk_action_execution_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE INDEX idx_action_execution_instance ON ordo_action_execution (instance_id); diff --git a/ordo-storage-jdbc/src/main/resources/db/postgresql/migration/V1__baseline.sql b/ordo-storage-jdbc/src/main/resources/db/postgresql/migration/V1__baseline.sql new file mode 100644 index 0000000..abf5823 --- /dev/null +++ b/ordo-storage-jdbc/src/main/resources/db/postgresql/migration/V1__baseline.sql @@ -0,0 +1,120 @@ +-- Ordo schema baseline (PostgreSQL / H2 MODE=PostgreSQL). + +CREATE TABLE ordo_process ( + id VARCHAR(64) PRIMARY KEY, + current_version INTEGER NOT NULL, + name VARCHAR(255) NOT NULL +); + +CREATE TABLE ordo_process_definition ( + id VARCHAR(64) NOT NULL, + version INTEGER NOT NULL, + name VARCHAR(255) NOT NULL, + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (id, version) +); + +CREATE TABLE ordo_approval_step ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + step_id VARCHAR(64) NOT NULL, + step_name VARCHAR(255) NOT NULL, + step_order INTEGER NOT NULL, + policy VARCHAR(16) NOT NULL, + kind VARCHAR(16) NOT NULL, + action_key VARCHAR(255), + due_after VARCHAR(32), + due_then VARCHAR(16), + due_to VARCHAR(255), + due_action VARCHAR(255), + PRIMARY KEY (definition_id, definition_version, step_id), + CONSTRAINT fk_approval_step_definition FOREIGN KEY (definition_id, definition_version) + REFERENCES ordo_process_definition (id, version) +); + +CREATE TABLE ordo_step_candidate ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + step_id VARCHAR(64) NOT NULL, + candidate VARCHAR(255) NOT NULL, + candidate_order INTEGER NOT NULL, + PRIMARY KEY (definition_id, definition_version, step_id, candidate), + CONSTRAINT fk_step_candidate_step FOREIGN KEY (definition_id, definition_version, step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id) +); + +CREATE TABLE ordo_step_transition ( + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + from_step_id VARCHAR(64) NOT NULL, + to_step_id VARCHAR(64), + condition_key VARCHAR(255), + priority INTEGER NOT NULL, + PRIMARY KEY (definition_id, definition_version, from_step_id, priority), + CONSTRAINT fk_transition_from FOREIGN KEY (definition_id, definition_version, from_step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id), + CONSTRAINT fk_transition_to FOREIGN KEY (definition_id, definition_version, to_step_id) + REFERENCES ordo_approval_step (definition_id, definition_version, step_id) +); + +CREATE TABLE ordo_process_instance ( + id VARCHAR(36) PRIMARY KEY, + definition_id VARCHAR(64) NOT NULL, + definition_version INTEGER NOT NULL, + initiator VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + context_json TEXT NOT NULL, + started_at TIMESTAMP NOT NULL, + finished_at TIMESTAMP, + CONSTRAINT fk_process_instance_definition FOREIGN KEY (definition_id, definition_version) + REFERENCES ordo_process_definition (id, version) +); + +CREATE TABLE ordo_approval_task ( + id VARCHAR(36) PRIMARY KEY, + instance_id VARCHAR(36) NOT NULL, + step_id VARCHAR(64) NOT NULL, + task_name VARCHAR(255) NOT NULL, + assignee VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + created_at TIMESTAMP NOT NULL, + completed_at TIMESTAMP, + action_actor VARCHAR(255), + action_comment TEXT, + action_at TIMESTAMP, + due_at TIMESTAMP, + CONSTRAINT fk_approval_task_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +); + +CREATE INDEX idx_approval_task_instance ON ordo_approval_task (instance_id); +CREATE INDEX idx_approval_task_status_assignee ON ordo_approval_task (status, assignee); +CREATE INDEX idx_approval_task_status_due ON ordo_approval_task (status, due_at); +CREATE INDEX idx_process_instance_definition ON ordo_process_instance (definition_id); + +CREATE TABLE ordo_process_event ( + id VARCHAR(36) PRIMARY KEY, + instance_id VARCHAR(36) NOT NULL, + task_id VARCHAR(36), + step_id VARCHAR(64), + event_type VARCHAR(32) NOT NULL, + actor VARCHAR(255), + detail TEXT, + occurred_at TIMESTAMP NOT NULL, + CONSTRAINT fk_process_event_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +); + +CREATE INDEX idx_process_event_instance_time ON ordo_process_event (instance_id, occurred_at, id); + +CREATE TABLE ordo_action_execution ( + id VARCHAR(36) PRIMARY KEY, + instance_id VARCHAR(36) NOT NULL, + step_id VARCHAR(64) NOT NULL, + action_key VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + error_message TEXT, + started_at TIMESTAMP NOT NULL, + finished_at TIMESTAMP, + CONSTRAINT fk_action_execution_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +); + +CREATE INDEX idx_action_execution_instance ON ordo_action_execution (instance_id); diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcMysqlIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcMysqlIntegrationTest.java new file mode 100644 index 0000000..e2e63e4 --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcMysqlIntegrationTest.java @@ -0,0 +1,282 @@ +package com.jetlumen.ordo.storage.jdbc; + +import com.jetlumen.ordo.api.ActionHandler; +import com.jetlumen.ordo.api.ApprovalStep; +import com.jetlumen.ordo.api.ApprovalTask; +import com.jetlumen.ordo.api.AssigneeResolver; +import com.jetlumen.ordo.api.OrdoEngine; +import com.jetlumen.ordo.api.ProcessContext; +import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.ProcessStatus; +import com.jetlumen.ordo.api.RoutingCondition; +import com.jetlumen.ordo.api.TaskAction; +import com.jetlumen.ordo.api.TaskStatus; +import com.jetlumen.ordo.core.DefaultOrdoEngine; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect; +import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects; +import com.mysql.cj.jdbc.MysqlDataSource; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Assumptions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import javax.sql.DataSource; +import java.sql.Connection; +import java.sql.SQLException; +import java.sql.Statement; +import java.time.Clock; +import java.time.Instant; +import java.time.ZoneOffset; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Optional MySQL counterpart of {@link JdbcPostgresIntegrationTest}. + * + *

Configured with {@code -Dordo.test.mysql.*} or {@code ORDO_TEST_MYSQL_*}; + * defaults {@code localhost:3306}, user/password {@code root}. Skipped when + * MySQL is unreachable. + */ +class JdbcMysqlIntegrationTest { + private static final String DATABASE = "ordo_test"; + private static final Instant NOW = Instant.parse("2026-03-03T10:00:00Z"); + private static final SqlDialect DIALECT = SqlDialects.mysql(); + + private static DataSource dataSource; + private JdbcConnectionProvider connectionProvider; + + @BeforeAll + static void setUpDatabase() throws SQLException { + MysqlDataSource rootDataSource = newDataSource(""); + String reachabilityProblem = null; + try (Connection ignored = rootDataSource.getConnection()) { + // probe + } catch (SQLException e) { + reachabilityProblem = e.getMessage(); + } + Assumptions.assumeTrue(reachabilityProblem == null, + "MySQL not reachable at " + config("ordo.test.mysql.host", "ORDO_TEST_MYSQL_HOST", "localhost") + + ":" + config("ordo.test.mysql.port", "ORDO_TEST_MYSQL_PORT", "3306") + + " — skipping MySQL integration tests (" + reachabilityProblem + ")"); + + try (Connection connection = rootDataSource.getConnection(); Statement statement = connection.createStatement()) { + statement.execute("DROP DATABASE IF EXISTS " + DATABASE); + statement.execute("CREATE DATABASE " + DATABASE + " CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci"); + } + + MysqlDataSource schemaDataSource = newDataSource(DATABASE); + JdbcTestSupport.applySchema(schemaDataSource, JdbcTestSupport.MYSQL_BASELINE); + dataSource = schemaDataSource; + } + + @AfterAll + static void tearDownDatabase() { + if (dataSource == null) { + return; + } + MysqlDataSource rootDataSource = newDataSource(""); + try (Connection connection = rootDataSource.getConnection(); Statement statement = connection.createStatement()) { + statement.execute("DROP DATABASE IF EXISTS " + DATABASE); + } catch (SQLException e) { + // leftover ordo_test is dropped on the next run + } + } + + @BeforeEach + void setUp() { + connectionProvider = new JdbcConnectionProvider(dataSource); + } + + @Test + void engineCompletesASequentialApprovalProcessOverMysql() { + OrdoEngine engine = newEngine(AssigneeResolver.direct()); + engine.publish(ProcessDefinition.linear("leave-mysql", "Leave request", List.of( + ApprovalStep.single("manager", "Manager approval", "maria"), + ApprovalStep.single("hr", "HR approval", "henry")))); + + ProcessInstance instance = engine.start("leave-mysql", "alice", + new ProcessContext(Map.of("requestId", "LEAVE-2026-001", "days", 5))); + ApprovalTask managerTask = engine.findPendingTasksByInstanceId(instance.id()).get(0); + engine.approve(managerTask.id(), "maria", "ok"); + ApprovalTask hrTask = engine.findPendingTasksByInstanceId(instance.id()).get(0); + engine.approve(hrTask.id(), "henry", "ok"); + + ProcessInstance finished = engine.findInstance(instance.id()).orElseThrow(); + assertEquals(ProcessStatus.APPROVED, finished.status()); + assertEquals(NOW, finished.finishedAt()); + assertEquals("LEAVE-2026-001", finished.context().value("requestId").orElseThrow()); + } + + @Test + void publishIsIdempotentForTheSameGraphAndVersionsAChangedGraph() { + JdbcProcessDefinitionRepository repository = + new JdbcProcessDefinitionRepository(connectionProvider, DIALECT); + ProcessDefinition definition = ProcessDefinition.linear("leave-dup-mysql", "Leave request", + List.of(ApprovalStep.single("manager", "Manager approval", "maria"))); + + assertEquals(1, repository.publish(definition).version()); + assertEquals(1, repository.publish(definition).version()); + assertEquals(2, repository.publish(ProcessDefinition.linear("leave-dup-mysql", "Second attempt", + List.of(ApprovalStep.single("manager", "Manager approval", "maria")))).version()); + assertEquals("Second attempt", repository.findLatest("leave-dup-mysql").orElseThrow().name()); + assertEquals("Leave request", repository.find("leave-dup-mysql", 1).orElseThrow().name()); + } + + @Test + void rejectsDuplicateInstanceIdsViaTheDatabaseUniqueConstraint() { + new JdbcProcessDefinitionRepository(connectionProvider, DIALECT).publish(ProcessDefinition.linear( + "leave-dup-inst-mysql", "Leave request", + List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); + JdbcProcessInstanceRepository repository = + new JdbcProcessInstanceRepository(connectionProvider, DIALECT); + ProcessInstance instance = new ProcessInstance("inst-dup-mysql", "leave-dup-inst-mysql", 1, "alice", + ProcessStatus.RUNNING, NOW, null, ProcessContext.empty()); + repository.insert(instance); + + assertThrows(JdbcStorageException.class, () -> repository.insert(instance)); + } + + @Test + void storesTimestampsAtMillisecondPrecision() { + new JdbcProcessDefinitionRepository(connectionProvider, DIALECT).publish(ProcessDefinition.linear( + "leave-time-mysql", "Leave request", + List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); + JdbcProcessInstanceRepository repository = + new JdbcProcessInstanceRepository(connectionProvider, DIALECT); + + Instant milliAligned = Instant.parse("2026-01-15T09:00:00.123Z"); + repository.insert(new ProcessInstance("inst-time-1", "leave-time-mysql", 1, "alice", + ProcessStatus.RUNNING, milliAligned, null, ProcessContext.empty())); + assertEquals(milliAligned, repository.findById("inst-time-1").orElseThrow().startedAt()); + + Instant withMicros = Instant.parse("2026-01-15T09:00:00.123456Z"); + repository.insert(new ProcessInstance("inst-time-2", "leave-time-mysql", 1, "alice", + ProcessStatus.RUNNING, withMicros, null, ProcessContext.empty())); + assertEquals(milliAligned, repository.findById("inst-time-2").orElseThrow().startedAt()); + } + + @Test + void onlyOneOfTwoConcurrentCompletionsWinsOnMysql() throws Exception { + new JdbcProcessDefinitionRepository(connectionProvider, DIALECT).publish(ProcessDefinition.linear( + "leave-race-mysql", "Leave request", + List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); + JdbcProcessInstanceRepository instanceRepository = + new JdbcProcessInstanceRepository(connectionProvider, DIALECT); + instanceRepository.insert(new ProcessInstance("inst-race-mysql", "leave-race-mysql", 1, "alice", + ProcessStatus.RUNNING, NOW, null, ProcessContext.empty())); + + JdbcApprovalTaskRepository taskRepository = + new JdbcApprovalTaskRepository(connectionProvider, DIALECT); + ApprovalTask pending = new ApprovalTask("task-race-mysql", "inst-race-mysql", "manager", + "Manager approval", "maria", TaskStatus.PENDING, NOW, null, null, null); + taskRepository.save(pending); + Instant completedAt = NOW.plusSeconds(30); + + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(2); + AtomicInteger wins = new AtomicInteger(); + Thread first = new Thread(() -> attempt(start, done, wins, completedTask(pending, "first", completedAt), + taskRepository)); + Thread second = new Thread(() -> attempt(start, done, wins, completedTask(pending, "second", completedAt), + taskRepository)); + first.start(); + second.start(); + start.countDown(); + + assertTrue(done.await(10, TimeUnit.SECONDS)); + assertEquals(1, wins.get()); + ApprovalTask stored = taskRepository.findById("task-race-mysql").orElseThrow(); + assertEquals("maria", stored.action().actor()); + assertTrue(Set.of("first", "second").contains(stored.action().comment())); + } + + @Test + void rollsBackTheWholeApprovalWhenTheNextStepCannotBeCreated() { + OrdoEngine failingEngine = newEngine((candidate, step, context) -> { + if (step.id().equals("hr")) { + throw new IllegalStateException("no hr approval today"); + } + return candidate; + }); + failingEngine.publish(ProcessDefinition.linear("leave-rollback-mysql", "Leave request", List.of( + ApprovalStep.single("manager", "Manager approval", "maria"), + ApprovalStep.single("hr", "HR approval", "henry")))); + + ProcessInstance instance = failingEngine.start("leave-rollback-mysql", "alice"); + ApprovalTask managerTask = failingEngine.findPendingTasksByInstanceId(instance.id()).get(0); + assertThrows(IllegalStateException.class, () -> failingEngine.approve(managerTask.id(), "maria")); + + ApprovalTask storedTask = failingEngine.findTask(managerTask.id()).orElseThrow(); + assertEquals(TaskStatus.PENDING, storedTask.status()); + assertNull(storedTask.action()); + assertEquals(1, failingEngine.findTasks(instance.id()).size()); + assertEquals(ProcessStatus.RUNNING, failingEngine.findInstance(instance.id()).orElseThrow().status()); + } + + private OrdoEngine newEngine(AssigneeResolver assigneeResolver) { + return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver, RoutingCondition.always(), + ActionHandler.noop(), + new JdbcTransactionExecutor(connectionProvider), + new JdbcProcessDefinitionRepository(connectionProvider, DIALECT), + new JdbcProcessInstanceRepository(connectionProvider, DIALECT), + new JdbcApprovalTaskRepository(connectionProvider, DIALECT), + new JdbcProcessHistoryRepository(connectionProvider, DIALECT), + new JdbcActionExecutionRepository(connectionProvider, DIALECT), + List.of()); + } + + private static MysqlDataSource newDataSource(String database) { + MysqlDataSource mysqlDataSource = new MysqlDataSource(); + String host = config("ordo.test.mysql.host", "ORDO_TEST_MYSQL_HOST", "localhost"); + String port = config("ordo.test.mysql.port", "ORDO_TEST_MYSQL_PORT", "3306"); + String databasePath = database.isEmpty() ? "/" : "/" + database; + mysqlDataSource.setUrl("jdbc:mysql://" + host + ":" + port + databasePath + + "?allowPublicKeyRetrieval=true&sslMode=DISABLED&characterEncoding=utf8"); + mysqlDataSource.setUser(config("ordo.test.mysql.user", "ORDO_TEST_MYSQL_USER", "root")); + mysqlDataSource.setPassword(config("ordo.test.mysql.password", "ORDO_TEST_MYSQL_PASSWORD", "root")); + return mysqlDataSource; + } + + private static String config(String property, String env, String defaultValue) { + String fromProperty = System.getProperty(property); + if (fromProperty != null && !fromProperty.isBlank()) { + return fromProperty; + } + String fromEnv = System.getenv(env); + if (fromEnv != null && !fromEnv.isBlank()) { + return fromEnv; + } + return defaultValue; + } + + private static void attempt(CountDownLatch start, CountDownLatch done, AtomicInteger wins, + ApprovalTask completed, JdbcApprovalTaskRepository repository) { + try { + start.await(); + if (repository.completeIfPending(completed)) { + wins.incrementAndGet(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + } + + private static ApprovalTask completedTask(ApprovalTask pending, String comment, Instant completedAt) { + return new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(), pending.name(), + pending.assignee(), TaskStatus.APPROVED, pending.createdAt(), completedAt, + new TaskAction(pending.assignee(), comment, completedAt), pending.dueAt()); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java index 048843d..796702f 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java @@ -49,7 +49,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue; * ({@code -Dordo.test.pg.host=myhost}) or environment variables * ({@code ORDO_TEST_PG_HOST}); defaults assume {@code localhost:5432} with * database/user/password {@code postgres}. Each run creates a dedicated - * {@code ordo_test} schema, applies the V1 migration there and drops it + * {@code ordo_test} schema, applies the PostgreSQL baseline there and drops it * afterward, so the rest of the database stays untouched. The class is * skipped when the database cannot be reached. */ diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java index c068c05..1305527 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java @@ -13,16 +13,8 @@ import java.util.UUID; /** Creates isolated in-memory H2 databases with the Ordo schema applied. */ final class JdbcTestSupport { - private static final String[] MIGRATIONS = { - "/db/migration/V1__create_ordo_tables.sql", - "/db/migration/V2__add_step_candidates_and_policy.sql", - "/db/migration/V3__add_step_transitions.sql", - "/db/migration/V4__add_step_kind_and_action_key.sql", - "/db/migration/V5__add_process_event_and_action_execution.sql", - "/db/migration/V6__add_step_due.sql", - "/db/migration/V7__definition_versions.sql" - }; - private static final String[] SCHEMA_SQL = loadSchemas(); + static final String POSTGRES_BASELINE = "/db/postgresql/migration/V1__baseline.sql"; + static final String MYSQL_BASELINE = "/db/mysql/migration/V1__baseline.sql"; private JdbcTestSupport() { } @@ -36,14 +28,16 @@ final class JdbcTestSupport { return dataSource; } - /** Applies every migration script in order to an empty database, e.g. a PostgreSQL test container. */ static void applySchema(DataSource dataSource) { + applySchema(dataSource, POSTGRES_BASELINE); + } + + static void applySchema(DataSource dataSource, String resourcePath) { + String migration = loadSchema(resourcePath); try (Connection connection = dataSource.getConnection(); Statement statement = connection.createStatement()) { - for (String migration : SCHEMA_SQL) { - for (String sql : stripComments(migration).split(";")) { - if (!sql.isBlank()) { - statement.execute(sql); - } + for (String sql : stripComments(migration).split(";")) { + if (!sql.isBlank()) { + statement.execute(sql); } } } catch (SQLException e) { @@ -64,14 +58,6 @@ final class JdbcTestSupport { return result.toString(); } - private static String[] loadSchemas() { - String[] schemas = new String[MIGRATIONS.length]; - for (int i = 0; i < MIGRATIONS.length; i++) { - schemas[i] = loadSchema(MIGRATIONS[i]); - } - return schemas; - } - private static String loadSchema(String resourcePath) { try (InputStream input = JdbcTestSupport.class.getResourceAsStream(resourcePath)) { if (input == null) { @@ -83,4 +69,3 @@ final class JdbcTestSupport { } } } - diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialectsTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialectsTest.java new file mode 100644 index 0000000..ff6f2eb --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/dialect/SqlDialectsTest.java @@ -0,0 +1,68 @@ +package com.jetlumen.ordo.storage.jdbc.dialect; + +import org.h2.jdbcx.JdbcDataSource; +import org.junit.jupiter.api.Test; + +import java.sql.Connection; +import java.sql.SQLException; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class SqlDialectsTest { + + @Test + void loadRegistersBuiltInDialects() { + assertTrue(SqlDialects.load().stream().anyMatch(d -> d.id().equals("postgresql"))); + assertTrue(SqlDialects.load().stream().anyMatch(d -> d.id().equals("mysql"))); + } + + @Test + void resolveByIdIgnoresDataSource() { + assertEquals("mysql", SqlDialects.resolve(null, "MySQL").id()); + assertEquals("postgresql", SqlDialects.resolve(null, "postgresql").id()); + } + + @Test + void unknownIdListsKnownDialects() { + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, + () -> SqlDialects.resolve(null, "oracle")); + assertTrue(ex.getMessage().contains("oracle")); + assertTrue(ex.getMessage().contains("postgresql")); + assertTrue(ex.getMessage().contains("mysql")); + } + + @Test + void detectsH2PostgreSQLMode() throws SQLException { + JdbcDataSource dataSource = new JdbcDataSource(); + dataSource.setURL("jdbc:h2:mem:dialect_" + UUID.randomUUID() + + ";MODE=PostgreSQL;DATABASE_TO_LOWER=TRUE;DB_CLOSE_DELAY=-1"); + dataSource.setUser("sa"); + try (Connection ignored = dataSource.getConnection()) { + assertEquals("postgresql", SqlDialects.resolve(dataSource, null).id()); + } + } + + @Test + void postgresDuplicateKeyIsSqlState23505() { + PostgresSqlDialect dialect = new PostgresSqlDialect(); + assertTrue(dialect.isDuplicateKey(new SQLException("dup", "23505"))); + assertFalse(dialect.isDuplicateKey(new SQLException("dup", "23000", 1062))); + assertEquals("SELECT 1 LIMIT ? OFFSET ?", dialect.limit("SELECT 1")); + assertEquals("SELECT 1 LIMIT ?", dialect.limit("SELECT 1", false)); + assertEquals("SELECT 1 FOR UPDATE", dialect.forUpdate("SELECT 1")); + assertEquals("classpath:db/postgresql/migration", dialect.flywayLocations()[0]); + } + + @Test + void mysqlDuplicateKeyIsSqlState23000AndError1062() { + MysqlSqlDialect dialect = new MysqlSqlDialect(); + assertTrue(dialect.isDuplicateKey(new SQLException("dup", "23000", 1062))); + assertFalse(dialect.isDuplicateKey(new SQLException("fk", "23000", 1452))); + assertFalse(dialect.isDuplicateKey(new SQLException("dup", "23505"))); + assertEquals("classpath:db/mysql/migration", dialect.flywayLocations()[0]); + } +} diff --git a/pom.xml b/pom.xml index be03a59..9df1fc5 100644 --- a/pom.xml +++ b/pom.xml @@ -30,6 +30,7 @@ 2.18.2 2.3.232 42.7.5 + 8.4.0 @@ -54,6 +55,11 @@ postgresql ${postgresql.version} + + com.mysql + mysql-connector-j + ${mysql.version} +