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 <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-16 11:05:14 +08:00
co-authored by Cursor
parent 9bda7c417c
commit 7d17157633
36 changed files with 1123 additions and 304 deletions
+5
View File
@@ -32,6 +32,11 @@
<artifactId>postgresql</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter</artifactId>
@@ -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<ActionExecution> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(sql)) {
select.setString(1, instanceId);
@@ -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<ApprovalTask> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(sql)) {
int index = PageSupport.bindParams(select, where.params());
@@ -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<ProcessDefinition> 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<int[]> 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<Integer> 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());
}
}
@@ -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<ProcessEvent> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(sql)) {
select.setString(1, instanceId);
@@ -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<ProcessInstance> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(sql)) {
int index = PageSupport.bindParams(select, where.params());
@@ -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";
}
}
@@ -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"};
}
}
@@ -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));
}
}
}
@@ -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();
}
@@ -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<SqlDialect> load() {
List<SqlDialect> 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<SqlDialect> 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<SqlDialect> 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<SqlDialect> 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<SqlDialect> dialects) {
return dialects.stream().map(SqlDialect::id).collect(Collectors.joining(", "));
}
}
@@ -0,0 +1,2 @@
com.jetlumen.ordo.storage.jdbc.dialect.PostgresSqlDialect
com.jetlumen.ordo.storage.jdbc.dialect.MysqlSqlDialect
@@ -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);
@@ -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;
@@ -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)
);
@@ -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);
@@ -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);
@@ -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);
@@ -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);
@@ -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);
@@ -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);
@@ -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}.
*
* <p>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());
}
}
@@ -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.
*/
@@ -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 {
}
}
}
@@ -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]);
}
}