feat: initial approval workflow engine with in-memory and JDBC storage

Lightweight linear approval engine (v0.1):

- ordo-api: domain model, OrdoEngine port, repository SPI with
  conditional updates (insertIfAbsent, completeIfPending) and a
  TransactionExecutor port for atomic multi-step writes
- ordo-core: DefaultOrdoEngine running register/start/approve/reject
  inside a transaction boundary; in-memory engine and repositories
- ordo-storage-jdbc: thread-bound JDBC transactions, normalized V1
  schema migration, Jackson-based ProcessContext JSON codec
- tests: unit tests plus H2 integration tests; PostgreSQL integration
  tests run against a local instance via ordo.test.pg.* properties and
  skip when unreachable

Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
0264408
2026-09-08 16:56:04 +08:00
co-authored by Claude Code
commit 397a9c229b
58 changed files with 2779 additions and 0 deletions
+47
View File
@@ -0,0 +1,47 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.jetlumen</groupId>
<artifactId>ordo</artifactId>
<version>1.0-SNAPSHOT</version>
</parent>
<artifactId>ordo-storage-jdbc</artifactId>
<dependencies>
<dependency>
<groupId>com.jetlumen</groupId>
<artifactId>ordo-api</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>com.jetlumen</groupId>
<artifactId>ordo-core</artifactId>
<version>${project.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.h2database</groupId>
<artifactId>h2</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,122 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
/** JDBC implementation of the task storage port and its pending-task indexes. */
public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository {
private static final String TASK_COLUMNS =
"id, instance_id, step_id, task_name, assignee, status, created_at, completed_at, action_actor, action_comment, action_at";
private static final String INSERT_TASK =
"INSERT INTO ordo_approval_task (id, instance_id, step_id, task_name, assignee, status, created_at)"
+ " VALUES (?, ?, ?, ?, ?, ?, ?)";
private static final String COMPLETE_IF_PENDING =
"UPDATE ordo_approval_task SET status = ?, completed_at = ?, action_actor = ?, action_comment = ?, action_at = ?"
+ " WHERE id = ? AND status = 'PENDING'";
private static final String SELECT_TASK =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE id = ?";
private static final String SELECT_BY_INSTANCE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE instance_id = ? ORDER BY created_at, id";
private static final String SELECT_PENDING_BY_ASSIGNEE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND assignee = ?"
+ " ORDER BY created_at, id";
private static final String SELECT_PENDING_BY_INSTANCE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND instance_id = ?"
+ " ORDER BY created_at, id";
private final JdbcConnectionProvider connectionProvider;
public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
}
@Override
public void save(ApprovalTask task) {
Objects.requireNonNull(task, "task must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement insert = connection.prepareStatement(INSERT_TASK)) {
ApprovalTaskMapper.bindInsert(insert, task);
insert.executeUpdate();
} catch (SQLException e) {
throw new JdbcStorageException("failed to insert task: " + task.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public Optional<ApprovalTask> findById(String taskId) {
return findOne(SELECT_TASK, taskId);
}
@Override
public List<ApprovalTask> findByInstanceId(String instanceId) {
return findAll(SELECT_BY_INSTANCE, instanceId);
}
@Override
public List<ApprovalTask> findPendingByAssignee(String assignee) {
return findAll(SELECT_PENDING_BY_ASSIGNEE, assignee);
}
@Override
public List<ApprovalTask> findPendingByInstanceId(String instanceId) {
return findAll(SELECT_PENDING_BY_INSTANCE, instanceId);
}
@Override
public boolean completeIfPending(ApprovalTask completedTask) {
Objects.requireNonNull(completedTask, "completedTask must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(COMPLETE_IF_PENDING)) {
ApprovalTaskMapper.bindComplete(update, completedTask);
return update.executeUpdate() == 1;
} catch (SQLException e) {
throw new JdbcStorageException("failed to complete task: " + completedTask.id(), e);
} finally {
connectionProvider.close(connection);
}
}
private Optional<ApprovalTask> findOne(String sql, String parameter) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) {
select.setString(1, parameter);
try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next() ? Optional.of(ApprovalTaskMapper.read(resultSet)) : Optional.empty();
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to query task: " + parameter, e);
} finally {
connectionProvider.close(connection);
}
}
private List<ApprovalTask> findAll(String sql, String parameter) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) {
select.setString(1, parameter);
try (ResultSet resultSet = select.executeQuery()) {
List<ApprovalTask> tasks = new ArrayList<>();
while (resultSet.next()) {
tasks.add(ApprovalTaskMapper.read(resultSet));
}
return tasks;
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to query tasks: " + parameter, e);
} finally {
connectionProvider.close(connection);
}
}
}
@@ -0,0 +1,74 @@
package com.jetlumen.ordo.storage.jdbc;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.SQLException;
import java.util.Objects;
/**
* Hands out JDBC connections to the repositories. When a transaction is
* active on the current thread, every call receives that transaction's
* connection; otherwise a new auto-commit connection is opened per call.
* Only {@link JdbcTransactionExecutor} starts and finishes transactions.
*/
public final class JdbcConnectionProvider {
private final DataSource dataSource;
private final ThreadLocal<Connection> transactionConnection = new ThreadLocal<>();
public JdbcConnectionProvider(DataSource dataSource) {
this.dataSource = Objects.requireNonNull(dataSource, "dataSource must not be null");
}
/** Returns the transaction-bound connection of the current thread, or opens a new auto-commit connection. */
public Connection getConnection() {
Connection connection = transactionConnection.get();
if (connection != null) {
return connection;
}
try {
return dataSource.getConnection();
} catch (SQLException e) {
throw new JdbcStorageException("failed to open a JDBC connection", e);
}
}
/** Closes a connection obtained from {@link #getConnection()}; a no-op while it belongs to the active transaction. */
public void close(Connection connection) {
if (connection == transactionConnection.get()) {
return; // released by closeTransaction when the transaction ends
}
try {
connection.close();
} catch (SQLException e) {
throw new JdbcStorageException("failed to close a JDBC connection", e);
}
}
boolean isTransactionActive() {
return transactionConnection.get() != null;
}
Connection openTransaction() {
if (isTransactionActive()) {
throw new IllegalStateException("a transaction is already active on this thread");
}
Connection connection = getConnection();
try {
connection.setAutoCommit(false);
} catch (SQLException e) {
close(connection);
throw new JdbcStorageException("failed to start a JDBC transaction", e);
}
transactionConnection.set(connection);
return connection;
}
void closeTransaction(Connection connection) {
transactionConnection.remove();
try {
connection.close();
} catch (SQLException e) {
throw new JdbcStorageException("failed to close a JDBC transaction connection", e);
}
}
}
@@ -0,0 +1,102 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
/** JDBC implementation of the definition storage port; steps live in a separate table. */
public final class JdbcProcessDefinitionRepository implements ProcessDefinitionRepository {
private static final String INSERT_DEFINITION =
"INSERT INTO ordo_process_definition (id, name) VALUES (?, ?)";
private static final String INSERT_STEP =
"INSERT INTO ordo_approval_step (definition_id, step_id, step_name, assignee, step_order) VALUES (?, ?, ?, ?, ?)";
private static final String SELECT_DEFINITION =
"SELECT id, name FROM ordo_process_definition WHERE id = ?";
private static final String SELECT_STEPS =
"SELECT step_id, step_name, assignee FROM ordo_approval_step WHERE definition_id = ? ORDER BY step_order";
private final JdbcConnectionProvider connectionProvider;
public JdbcProcessDefinitionRepository(JdbcConnectionProvider connectionProvider) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
}
@Override
public boolean insertIfAbsent(ProcessDefinition definition) {
Objects.requireNonNull(definition, "definition must not be null");
Connection connection = connectionProvider.getConnection();
try {
try (PreparedStatement insertDefinition = connection.prepareStatement(INSERT_DEFINITION)) {
insertDefinition.setString(1, definition.id());
insertDefinition.setString(2, definition.name());
insertDefinition.executeUpdate();
} catch (SQLException e) {
if (isDuplicateKey(e)) {
return false;
}
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e);
}
int order = 0;
for (ApprovalStep step : definition.steps()) {
try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) {
insertStep.setString(1, definition.id());
insertStep.setString(2, step.id());
insertStep.setString(3, step.name());
insertStep.setString(4, step.assignee());
insertStep.setInt(5, order++);
insertStep.executeUpdate();
}
}
return true;
} catch (SQLException e) {
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public Optional<ProcessDefinition> findById(String definitionId) {
Connection connection = connectionProvider.getConnection();
try {
String name;
try (PreparedStatement selectDefinition = connection.prepareStatement(SELECT_DEFINITION)) {
selectDefinition.setString(1, definitionId);
try (ResultSet resultSet = selectDefinition.executeQuery()) {
if (!resultSet.next()) {
return Optional.empty();
}
name = resultSet.getString("name");
}
}
List<ApprovalStep> steps = new ArrayList<>();
try (PreparedStatement selectSteps = connection.prepareStatement(SELECT_STEPS)) {
selectSteps.setString(1, definitionId);
try (ResultSet resultSet = selectSteps.executeQuery()) {
while (resultSet.next()) {
steps.add(ApprovalStepMapper.read(resultSet));
}
}
}
return Optional.of(new ProcessDefinition(definitionId, name, steps));
} catch (SQLException e) {
throw new JdbcStorageException("failed to load definition: " + definitionId, e);
} finally {
connectionProvider.close(connection);
}
}
private static boolean isDuplicateKey(SQLException e) {
return "23505".equals(e.getSQLState());
}
}
@@ -0,0 +1,75 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.Objects;
import java.util.Optional;
/** JDBC implementation of the instance storage port. */
public final class JdbcProcessInstanceRepository implements ProcessInstanceRepository {
private static final String INSERT_INSTANCE =
"INSERT INTO ordo_process_instance (id, definition_id, initiator, status, context_json, started_at, finished_at)"
+ " VALUES (?, ?, ?, ?, ?, ?, ?)";
private static final String UPDATE_INSTANCE =
"UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ?";
private static final String SELECT_INSTANCE =
"SELECT id, definition_id, initiator, status, context_json, started_at, finished_at"
+ " FROM ordo_process_instance WHERE id = ?";
private final JdbcConnectionProvider connectionProvider;
public JdbcProcessInstanceRepository(JdbcConnectionProvider connectionProvider) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
}
@Override
public void insert(ProcessInstance instance) {
Objects.requireNonNull(instance, "instance must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement insert = connection.prepareStatement(INSERT_INSTANCE)) {
ProcessInstanceMapper.bindInsert(insert, instance);
insert.executeUpdate();
} catch (SQLException e) {
throw new JdbcStorageException("failed to insert instance: " + instance.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public void update(ProcessInstance instance) {
Objects.requireNonNull(instance, "instance must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(UPDATE_INSTANCE)) {
ProcessInstanceMapper.bindUpdate(update, instance);
if (update.executeUpdate() != 1) {
throw new IllegalStateException("instance not found: " + instance.id());
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to update instance: " + instance.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public Optional<ProcessInstance> findById(String instanceId) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(SELECT_INSTANCE)) {
select.setString(1, instanceId);
try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next() ? Optional.of(ProcessInstanceMapper.read(resultSet)) : Optional.empty();
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to load instance: " + instanceId, e);
} finally {
connectionProvider.close(connection);
}
}
}
@@ -0,0 +1,10 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.exception.OrdoException;
/** Unchecked wrapper for JDBC failures raised by the JDBC storage module. */
public final class JdbcStorageException extends OrdoException {
public JdbcStorageException(String message, Throwable cause) {
super(message, cause);
}
}
@@ -0,0 +1,57 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.TransactionExecutor;
import java.sql.Connection;
import java.sql.SQLException;
import java.util.Objects;
import java.util.function.Supplier;
/**
* Runs actions inside a JDBC transaction bound to the current thread. Every
* repository operation executed by the action shares the same connection and
* is committed or rolled back together. A nested {@code execute} joins the
* surrounding transaction.
*/
public final class JdbcTransactionExecutor implements TransactionExecutor {
private final JdbcConnectionProvider connectionProvider;
public JdbcTransactionExecutor(JdbcConnectionProvider connectionProvider) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
}
@Override
public <T> T execute(Supplier<T> action) {
Objects.requireNonNull(action, "action must not be null");
if (connectionProvider.isTransactionActive()) {
return action.get(); // nested execution joins the surrounding transaction
}
Connection connection = connectionProvider.openTransaction();
try {
T result = action.get();
commit(connection);
return result;
} catch (RuntimeException | Error failure) {
rollback(connection, failure);
throw failure;
} finally {
connectionProvider.closeTransaction(connection);
}
}
private static void commit(Connection connection) {
try {
connection.commit();
} catch (SQLException e) {
throw new JdbcStorageException("failed to commit the JDBC transaction", e);
}
}
private static void rollback(Connection connection, Throwable failure) {
try {
connection.rollback();
} catch (SQLException e) {
failure.addSuppressed(e);
}
}
}
@@ -0,0 +1,17 @@
package com.jetlumen.ordo.storage.jdbc.mapper;
import com.jetlumen.ordo.api.ApprovalStep;
import java.sql.ResultSet;
import java.sql.SQLException;
/** Maps rows of {@code ordo_approval_step} to {@link ApprovalStep} objects. */
public final class ApprovalStepMapper {
private ApprovalStepMapper() {
}
public static ApprovalStep read(ResultSet resultSet) throws SQLException {
return new ApprovalStep(resultSet.getString("step_id"), resultSet.getString("step_name"),
resultSet.getString("assignee"));
}
}
@@ -0,0 +1,54 @@
package com.jetlumen.ordo.storage.jdbc.mapper;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.TaskAction;
import com.jetlumen.ordo.api.TaskStatus;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
/** Maps rows of {@code ordo_approval_task} to {@link ApprovalTask} objects and back. */
public final class ApprovalTaskMapper {
private ApprovalTaskMapper() {
}
public static void bindInsert(PreparedStatement statement, ApprovalTask task) throws SQLException {
statement.setString(1, task.id());
statement.setString(2, task.instanceId());
statement.setString(3, task.stepId());
statement.setString(4, task.name());
statement.setString(5, task.assignee());
statement.setString(6, task.status().name());
statement.setTimestamp(7, Timestamp.from(task.createdAt()));
}
public static void bindComplete(PreparedStatement statement, ApprovalTask completedTask) throws SQLException {
TaskAction action = completedTask.action();
statement.setString(1, completedTask.status().name());
statement.setTimestamp(2, Timestamp.from(completedTask.completedAt()));
statement.setString(3, action.actor());
statement.setString(4, action.comment());
statement.setTimestamp(5, Timestamp.from(action.operatedAt()));
statement.setString(6, completedTask.id());
}
public static ApprovalTask read(ResultSet resultSet) throws SQLException {
String actionActor = resultSet.getString("action_actor");
Timestamp actionAt = resultSet.getTimestamp("action_at");
TaskAction action = actionActor == null ? null
: new TaskAction(actionActor, resultSet.getString("action_comment"), actionAt.toInstant());
Timestamp completedAt = resultSet.getTimestamp("completed_at");
return new ApprovalTask(
resultSet.getString("id"),
resultSet.getString("instance_id"),
resultSet.getString("step_id"),
resultSet.getString("task_name"),
resultSet.getString("assignee"),
TaskStatus.valueOf(resultSet.getString("status")),
resultSet.getTimestamp("created_at").toInstant(),
completedAt == null ? null : completedAt.toInstant(),
action);
}
}
@@ -0,0 +1,38 @@
package com.jetlumen.ordo.storage.jdbc.mapper;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.jetlumen.ordo.api.ProcessContext;
import java.util.Map;
/**
* Serializes a {@link ProcessContext} to and from the {@code context_json}
* column. v0.1 only supports JSON-compatible values: strings, numbers,
* booleans, lists and nested maps.
*/
public final class ProcessContextCodec {
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
private ProcessContextCodec() {
}
public static String encode(ProcessContext context) {
try {
return OBJECT_MAPPER.writeValueAsString(context.variables());
} catch (JsonProcessingException e) {
throw new IllegalStateException("failed to serialize process context to JSON", e);
}
}
public static ProcessContext decode(String json) {
try {
Map<String, Object> variables = OBJECT_MAPPER.readValue(json, new TypeReference<>() {
});
return new ProcessContext(variables);
} catch (JsonProcessingException e) {
throw new IllegalStateException("failed to deserialize process context from JSON", e);
}
}
}
@@ -0,0 +1,44 @@
package com.jetlumen.ordo.storage.jdbc.mapper;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
/** Maps rows of {@code ordo_process_instance} to {@link ProcessInstance} objects and back. */
public final class ProcessInstanceMapper {
private ProcessInstanceMapper() {
}
public static void bindInsert(PreparedStatement statement, ProcessInstance instance) throws SQLException {
statement.setString(1, instance.id());
statement.setString(2, instance.definitionId());
statement.setString(3, instance.initiator());
statement.setString(4, instance.status().name());
statement.setString(5, ProcessContextCodec.encode(instance.context()));
statement.setTimestamp(6, Timestamp.from(instance.startedAt()));
statement.setTimestamp(7, instance.finishedAt() == null ? null : Timestamp.from(instance.finishedAt()));
}
public static void bindUpdate(PreparedStatement statement, ProcessInstance instance) throws SQLException {
statement.setString(1, instance.status().name());
statement.setTimestamp(2, instance.finishedAt() == null ? null : Timestamp.from(instance.finishedAt()));
statement.setString(3, instance.id());
}
public static ProcessInstance read(ResultSet resultSet) throws SQLException {
Timestamp startedAt = resultSet.getTimestamp("started_at");
Timestamp finishedAt = resultSet.getTimestamp("finished_at");
return new ProcessInstance(
resultSet.getString("id"),
resultSet.getString("definition_id"),
resultSet.getString("initiator"),
ProcessStatus.valueOf(resultSet.getString("status")),
startedAt.toInstant(),
finishedAt == null ? null : finishedAt.toInstant(),
ProcessContextCodec.decode(resultSet.getString("context_json")));
}
}
@@ -0,0 +1,48 @@
-- 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);
@@ -0,0 +1,143 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
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.TaskAction;
import com.jetlumen.ordo.api.TaskStatus;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.time.Instant;
import java.util.List;
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.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
class JdbcApprovalTaskRepositoryTest {
private static final Instant CREATED_AT = Instant.parse("2026-01-15T09:00:00Z");
private static final Instant COMPLETED_AT = CREATED_AT.plusSeconds(30);
private JdbcConnectionProvider connectionProvider;
private JdbcApprovalTaskRepository repository;
@BeforeEach
void setUp() {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
repository = new JdbcApprovalTaskRepository(connectionProvider);
insertFixtureData();
}
/** Tasks reference their instance, which references its definition; both parent rows must exist. */
private void insertFixtureData() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(new ProcessDefinition("leave",
"Leave request", List.of(new ApprovalStep("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider);
instanceRepository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty()));
instanceRepository.insert(new ProcessInstance("inst-2", "leave", "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty()));
}
@Test
void roundTripsAPendingTaskAndItsCompletedState() {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
repository.save(pending);
assertEquals(pending, repository.findById("task-1").orElseThrow());
ApprovalTask completed = completedTask(pending, "ok");
assertTrue(repository.completeIfPending(completed));
assertEquals(completed, repository.findById("task-1").orElseThrow());
}
@Test
void completeIfPendingFailsForAnAlreadyCompletedTask() {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
repository.save(pending);
assertTrue(repository.completeIfPending(completedTask(pending, "first")));
assertFalse(repository.completeIfPending(completedTask(pending, "second")));
assertEquals("first", repository.findById("task-1").orElseThrow().action().comment());
}
@Test
void completeIfPendingFailsForAnUnknownTask() {
ApprovalTask pending = pendingTask("missing", "inst-1", "manager", "maria", CREATED_AT);
assertFalse(repository.completeIfPending(completedTask(pending, "ok")));
}
@Test
void findsPendingTasksByAssigneeAndInstance() {
ApprovalTask mariaTask = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
ApprovalTask henryTask = pendingTask("task-2", "inst-1", "hr", "henry", CREATED_AT.plusSeconds(5));
ApprovalTask otherMariaTask = pendingTask("task-3", "inst-2", "manager", "maria", CREATED_AT.plusSeconds(10));
repository.save(mariaTask);
repository.save(henryTask);
repository.save(otherMariaTask);
assertEquals(List.of(mariaTask, otherMariaTask), repository.findPendingByAssignee("maria"));
assertEquals(List.of(mariaTask, henryTask), repository.findPendingByInstanceId("inst-1"));
assertEquals(List.of(mariaTask, henryTask), repository.findByInstanceId("inst-1"));
assertEquals(List.of(otherMariaTask), repository.findByInstanceId("inst-2"));
repository.completeIfPending(completedTask(mariaTask, "ok"));
assertEquals(List.of(otherMariaTask), repository.findPendingByAssignee("maria"));
assertEquals(List.of(henryTask), repository.findPendingByInstanceId("inst-1"));
assertEquals(List.of(henryTask), repository.findPendingByAssignee("henry"));
}
@Test
void onlyOneOfTwoConcurrentCompletionsWins() throws Exception {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
repository.save(pending);
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")));
Thread second = new Thread(() -> attempt(start, done, wins, completedTask(pending, "second")));
first.start();
second.start();
start.countDown();
assertTrue(done.await(10, TimeUnit.SECONDS));
assertEquals(1, wins.get());
ApprovalTask stored = repository.findById("task-1").orElseThrow();
assertEquals("maria", stored.action().actor());
assertTrue(Set.of("first", "second").contains(stored.action().comment()));
}
private void attempt(CountDownLatch start, CountDownLatch done, AtomicInteger wins, ApprovalTask completed) {
try {
start.await();
if (repository.completeIfPending(completed)) {
wins.incrementAndGet();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
}
private static ApprovalTask pendingTask(String id, String instanceId, String stepId, String assignee,
Instant createdAt) {
return new ApprovalTask(id, instanceId, stepId, stepId + " approval", assignee,
TaskStatus.PENDING, createdAt, null, null);
}
private static ApprovalTask completedTask(ApprovalTask pending, String comment) {
return new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(), pending.name(),
pending.assignee(), TaskStatus.APPROVED, pending.createdAt(), COMPLETED_AT,
new TaskAction(pending.assignee(), comment, COMPLETED_AT));
}
}
@@ -0,0 +1,133 @@
package com.jetlumen.ordo.storage.jdbc;
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.TaskStatus;
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
import com.jetlumen.ordo.core.DefaultOrdoEngine;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.time.Clock;
import java.time.Instant;
import java.time.ZoneOffset;
import java.util.List;
import java.util.Map;
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;
class JdbcOrdoEngineIntegrationTest {
private static final Instant NOW = Instant.parse("2026-02-02T10:00:00Z");
private JdbcConnectionProvider connectionProvider;
private OrdoEngine engine;
@BeforeEach
void setUp() {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
engine = newEngine(AssigneeResolver.direct());
engine.register(new ProcessDefinition("leave", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry"))));
}
@Test
void completesASequentialApprovalProcess() {
ProcessInstance instance = engine.start("leave", "alice",
new ProcessContext(Map.of("requestId", "LEAVE-2026-001")));
ApprovalTask managerTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertEquals("maria", managerTask.assignee());
engine.approve(managerTask.id(), "maria", "ok");
ApprovalTask hrTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertEquals("henry", hrTask.assignee());
assertEquals(ProcessStatus.RUNNING, engine.findInstance(instance.id()).orElseThrow().status());
engine.approve(hrTask.id(), "henry", "ok");
ProcessInstance finished = engine.findInstance(instance.id()).orElseThrow();
assertEquals(ProcessStatus.APPROVED, finished.status());
assertEquals(NOW, finished.finishedAt());
assertTrue(engine.findPendingTasksByInstanceId(instance.id()).isEmpty());
assertEquals("LEAVE-2026-001", finished.context().value("requestId").orElseThrow());
}
@Test
void rejectionTerminatesTheProcess() {
ProcessInstance instance = engine.start("leave", "alice");
ApprovalTask task = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
ApprovalTask rejectedTask = engine.reject(task.id(), "maria", "Insufficient leave balance");
assertEquals(TaskStatus.REJECTED, rejectedTask.status());
assertEquals(ProcessStatus.REJECTED, engine.findInstance(instance.id()).orElseThrow().status());
assertEquals(1, engine.findTasks(instance.id()).size());
assertThrows(TaskAlreadyCompletedException.class, () -> engine.approve(task.id(), "maria"));
}
@Test
void rejectsRepeatedApprovalsWithoutDuplicatingTheFlow() {
ProcessInstance instance = engine.start("leave", "alice");
ApprovalTask managerTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
engine.approve(managerTask.id(), "maria");
assertThrows(TaskAlreadyCompletedException.class, () -> engine.approve(managerTask.id(), "maria"));
assertEquals(2, engine.findTasks(instance.id()).size());
assertEquals(1, engine.findPendingTasksByInstanceId(instance.id()).size());
assertEquals("hr", engine.findPendingTasksByInstanceId(instance.id()).getFirst().stepId());
}
@Test
void rollsBackTheWholeApprovalWhenTheNextStepCannotBeCreated() {
OrdoEngine failingEngine = newEngine((step, context) -> {
if (step.id().equals("hr")) {
throw new IllegalStateException("no hr approval today");
}
return step.assignee();
});
failingEngine.register(new ProcessDefinition("leave2", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry"))));
ProcessInstance instance = failingEngine.start("leave2", "alice");
ApprovalTask managerTask = failingEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertThrows(IllegalStateException.class, () -> failingEngine.approve(managerTask.id(), "maria"));
// the conditional task completion must have been rolled back
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());
}
@Test
void dataSurvivesAcrossEngineInstancesOverTheSameDataSource() {
ProcessInstance instance = engine.start("leave", "alice");
OrdoEngine secondEngine = newEngine(AssigneeResolver.direct());
assertEquals(instance, secondEngine.findInstance(instance.id()).orElseThrow());
ApprovalTask managerTask = secondEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
secondEngine.approve(managerTask.id(), "maria");
assertEquals(ProcessStatus.RUNNING, engine.findInstance(instance.id()).orElseThrow().status());
assertEquals("henry", engine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee());
}
private OrdoEngine newEngine(AssigneeResolver assigneeResolver) {
return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver,
new JdbcTransactionExecutor(connectionProvider),
new JdbcProcessDefinitionRepository(connectionProvider),
new JdbcProcessInstanceRepository(connectionProvider),
new JdbcApprovalTaskRepository(connectionProvider));
}
}
@@ -0,0 +1,271 @@
package com.jetlumen.ordo.storage.jdbc;
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.TaskAction;
import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.core.DefaultOrdoEngine;
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 org.postgresql.ds.PGSimpleDataSource;
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.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Verifies the JDBC module against a local PostgreSQL instance, complementing
* the H2 tests: concurrency of conditional updates, TIMESTAMP semantics and
* unique constraints behave the same as in production.
*
* <p>The target database is configured with system properties
* ({@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
* afterward, so the rest of the database stays untouched. The class is
* skipped when the database cannot be reached.
*/
class JdbcPostgresIntegrationTest {
private static final String SCHEMA = "ordo_test";
private static final Instant NOW = Instant.parse("2026-03-03T10:00:00Z");
private static DataSource dataSource;
private JdbcConnectionProvider connectionProvider;
@BeforeAll
static void setUpDatabase() throws SQLException {
PGSimpleDataSource rootDataSource = newDataSource("");
String reachabilityProblem = null;
try (Connection ignored = rootDataSource.getConnection()) {
// just probing reachability
} catch (SQLException e) {
reachabilityProblem = e.getMessage();
}
Assumptions.assumeTrue(reachabilityProblem == null,
"PostgreSQL not reachable at " + rootDataSource.getServerNames()[0] + ":"
+ rootDataSource.getPortNumbers()[0] + "/" + rootDataSource.getDatabaseName()
+ " — skipping PG integration tests (" + reachabilityProblem + ")");
try (Connection connection = rootDataSource.getConnection(); Statement statement = connection.createStatement()) {
statement.execute("DROP SCHEMA IF EXISTS " + SCHEMA + " CASCADE");
statement.execute("CREATE SCHEMA " + SCHEMA);
}
PGSimpleDataSource schemaDataSource = newDataSource(SCHEMA);
JdbcTestSupport.applySchema(schemaDataSource);
dataSource = schemaDataSource;
}
@AfterAll
static void tearDownDatabase() {
if (dataSource == null) {
return;
}
try (Connection connection = dataSource.getConnection(); Statement statement = connection.createStatement()) {
statement.execute("DROP SCHEMA IF EXISTS " + SCHEMA + " CASCADE");
} catch (SQLException e) {
// a leftover ordo_test schema is harmless; the next run drops it again
}
}
@BeforeEach
void setUp() {
connectionProvider = new JdbcConnectionProvider(dataSource);
}
@Test
void engineCompletesASequentialApprovalProcessOverPostgres() {
OrdoEngine engine = newEngine(AssigneeResolver.direct());
engine.register(new ProcessDefinition("leave-pg", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry"))));
ProcessInstance instance = engine.start("leave-pg", "alice",
new ProcessContext(Map.of("requestId", "LEAVE-2026-001", "days", 5)));
ApprovalTask managerTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
engine.approve(managerTask.id(), "maria", "ok");
ApprovalTask hrTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst();
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 rejectsDuplicateDefinitionIdsViaTheDatabaseUniqueConstraint() {
JdbcProcessDefinitionRepository repository = new JdbcProcessDefinitionRepository(connectionProvider);
ProcessDefinition definition = new ProcessDefinition("leave-dup-pg", "Leave request",
List.of(new ApprovalStep("manager", "Manager approval", "maria")));
assertTrue(repository.insertIfAbsent(definition));
assertFalse(repository.insertIfAbsent(new ProcessDefinition("leave-dup-pg", "Second attempt",
List.of(new ApprovalStep("manager", "Manager approval", "maria")))));
assertEquals("Leave request", repository.findById("leave-dup-pg").orElseThrow().name());
}
@Test
void rejectsDuplicateInstanceIdsViaTheDatabaseUniqueConstraint() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(new ProcessDefinition(
"leave-dup-inst-pg", "Leave request", List.of(new ApprovalStep("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider);
ProcessInstance instance = new ProcessInstance("inst-dup-pg", "leave-dup-inst-pg", "alice",
ProcessStatus.RUNNING, NOW, null, ProcessContext.empty());
repository.insert(instance);
assertThrows(JdbcStorageException.class, () -> repository.insert(instance));
}
@Test
void roundsTimestampsToMicrosecondPrecision() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(new ProcessDefinition(
"leave-time-pg", "Leave request", List.of(new ApprovalStep("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider);
Instant microAligned = Instant.parse("2026-01-15T09:00:00.123456Z");
repository.insert(new ProcessInstance("inst-time-1", "leave-time-pg", "alice",
ProcessStatus.RUNNING, microAligned, null, ProcessContext.empty()));
assertEquals(microAligned, repository.findById("inst-time-1").orElseThrow().startedAt());
// PostgreSQL's TIMESTAMP stores microseconds and rounds the fractional seconds
Instant withNanos = Instant.parse("2026-01-15T09:00:00.123456789Z");
repository.insert(new ProcessInstance("inst-time-2", "leave-time-pg", "alice",
ProcessStatus.RUNNING, withNanos, null, ProcessContext.empty()));
assertEquals(Instant.parse("2026-01-15T09:00:00.123457Z"),
repository.findById("inst-time-2").orElseThrow().startedAt());
}
@Test
void onlyOneOfTwoConcurrentCompletionsWinsOnPostgres() throws Exception {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(new ProcessDefinition(
"leave-race-pg", "Leave request", List.of(new ApprovalStep("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider);
instanceRepository.insert(new ProcessInstance("inst-race-pg", "leave-race-pg", "alice",
ProcessStatus.RUNNING, NOW, null, ProcessContext.empty()));
JdbcApprovalTaskRepository taskRepository = new JdbcApprovalTaskRepository(connectionProvider);
ApprovalTask pending = new ApprovalTask("task-race-pg", "inst-race-pg", "manager", "Manager approval",
"maria", TaskStatus.PENDING, NOW, 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-pg").orElseThrow();
assertEquals("maria", stored.action().actor());
assertTrue(Set.of("first", "second").contains(stored.action().comment()));
}
@Test
void rollsBackTheWholeApprovalWhenTheNextStepCannotBeCreated() {
OrdoEngine failingEngine = newEngine((step, context) -> {
if (step.id().equals("hr")) {
throw new IllegalStateException("no hr approval today");
}
return step.assignee();
});
failingEngine.register(new ProcessDefinition("leave-rollback-pg", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry"))));
ProcessInstance instance = failingEngine.start("leave-rollback-pg", "alice");
ApprovalTask managerTask = failingEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertThrows(IllegalStateException.class, () -> failingEngine.approve(managerTask.id(), "maria"));
// the conditional task completion must have been rolled back
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,
new JdbcTransactionExecutor(connectionProvider),
new JdbcProcessDefinitionRepository(connectionProvider),
new JdbcProcessInstanceRepository(connectionProvider),
new JdbcApprovalTaskRepository(connectionProvider));
}
private static PGSimpleDataSource newDataSource(String currentSchema) {
PGSimpleDataSource pgDataSource = new PGSimpleDataSource();
pgDataSource.setServerNames(new String[]{config("ordo.test.pg.host", "ORDO_TEST_PG_HOST", "localhost")});
pgDataSource.setPortNumbers(new int[]{Integer.parseInt(config("ordo.test.pg.port", "ORDO_TEST_PG_PORT", "5432"))});
pgDataSource.setDatabaseName(config("ordo.test.pg.database", "ORDO_TEST_PG_DATABASE", "postgres"));
pgDataSource.setUser(config("ordo.test.pg.user", "ORDO_TEST_PG_USER", "postgres"));
pgDataSource.setPassword(config("ordo.test.pg.password", "ORDO_TEST_PG_PASSWORD", "postgres"));
if (!currentSchema.isEmpty()) {
pgDataSource.setCurrentSchema(currentSchema);
}
return pgDataSource;
}
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));
}
}
@@ -0,0 +1,49 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ProcessDefinition;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
class JdbcProcessDefinitionRepositoryTest {
private JdbcProcessDefinitionRepository repository;
@BeforeEach
void setUp() {
repository = new JdbcProcessDefinitionRepository(
new JdbcConnectionProvider(JdbcTestSupport.newDataSource()));
}
@Test
void insertsAndReadsBackADefinitionWithItsStepsInOrder() {
ProcessDefinition definition = new ProcessDefinition("leave", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry")));
assertTrue(repository.insertIfAbsent(definition));
assertEquals(definition, repository.findById("leave").orElseThrow());
}
@Test
void rejectsAnExistingDefinitionId() {
assertTrue(repository.insertIfAbsent(definition("leave", "Leave request v1")));
assertFalse(repository.insertIfAbsent(definition("leave", "Leave request v2")));
assertEquals("Leave request v1", repository.findById("leave").orElseThrow().name());
}
@Test
void returnsEmptyForAnUnknownDefinition() {
assertTrue(repository.findById("missing").isEmpty());
}
private static ProcessDefinition definition(String id, String name) {
return new ProcessDefinition(id, name, List.of(new ApprovalStep("lead", "Lead approval", "lee")));
}
}
@@ -0,0 +1,74 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep;
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 org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
class JdbcProcessInstanceRepositoryTest {
private static final Instant STARTED_AT = Instant.parse("2026-01-15T09:00:00Z");
private JdbcConnectionProvider connectionProvider;
private JdbcProcessInstanceRepository repository;
@BeforeEach
void setUp() {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
repository = new JdbcProcessInstanceRepository(connectionProvider);
// instances reference their definition, so the parent row must exist
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(new ProcessDefinition("leave",
"Leave request", List.of(new ApprovalStep("manager", "Manager approval", "maria"))));
}
@Test
void roundTripsAnInstanceIncludingItsJsonContext() {
ProcessContext context = new ProcessContext(Map.of(
"requestId", "LEAVE-2026-001",
"days", 5,
"urgent", true,
"candidates", List.of("maria", "henry"),
"meta", Map.of("priority", "high", "retries", 2)));
ProcessInstance instance = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING,
STARTED_AT, null, context);
repository.insert(instance);
assertEquals(instance, repository.findById("inst-1").orElseThrow());
}
@Test
void updatesStatusAndFinishedAt() {
repository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty()));
Instant finishedAt = STARTED_AT.plusSeconds(300);
repository.update(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.APPROVED,
STARTED_AT, finishedAt, ProcessContext.empty()));
ProcessInstance updated = repository.findById("inst-1").orElseThrow();
assertEquals(ProcessStatus.APPROVED, updated.status());
assertEquals(finishedAt, updated.finishedAt());
}
@Test
void returnsEmptyForAnUnknownInstance() {
assertTrue(repository.findById("missing").isEmpty());
}
@Test
void updateOfAnUnknownInstanceFails() {
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",
"alice", ProcessStatus.APPROVED, STARTED_AT, STARTED_AT, ProcessContext.empty())));
}
}
@@ -0,0 +1,53 @@
package com.jetlumen.ordo.storage.jdbc;
import org.h2.jdbcx.JdbcDataSource;
import javax.sql.DataSource;
import java.io.IOException;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.UUID;
/** Creates isolated in-memory H2 databases with the Ordo schema applied. */
final class JdbcTestSupport {
private static final String SCHEMA_SQL = loadSchema();
private JdbcTestSupport() {
}
static DataSource newDataSource() {
JdbcDataSource dataSource = new JdbcDataSource();
dataSource.setURL("jdbc:h2:mem:ordo_" + UUID.randomUUID()
+ ";MODE=PostgreSQL;DATABASE_TO_LOWER=TRUE;DB_CLOSE_DELAY=-1");
dataSource.setUser("sa");
applySchema(dataSource);
return dataSource;
}
/** Applies the V1 migration script to an empty database, e.g. a PostgreSQL test container. */
static void applySchema(DataSource dataSource) {
try (Connection connection = dataSource.getConnection(); Statement statement = connection.createStatement()) {
for (String sql : SCHEMA_SQL.split(";")) {
if (!sql.isBlank()) {
statement.execute(sql);
}
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to apply the Ordo schema", e);
}
}
private static String loadSchema() {
try (InputStream input = JdbcTestSupport.class.getResourceAsStream("/db/migration/V1__create_ordo_tables.sql")) {
if (input == null) {
throw new IllegalStateException("V1__create_ordo_tables.sql not found on the classpath");
}
return new String(input.readAllBytes(), StandardCharsets.UTF_8);
} catch (IOException e) {
throw new IllegalStateException("failed to read V1__create_ordo_tables.sql", e);
}
}
}
@@ -0,0 +1,88 @@
package com.jetlumen.ordo.storage.jdbc;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
class JdbcTransactionExecutorTest {
private JdbcConnectionProvider connectionProvider;
private JdbcTransactionExecutor transactionExecutor;
@BeforeEach
void setUp() {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
transactionExecutor = new JdbcTransactionExecutor(connectionProvider);
}
@Test
void commitsAllStatementsWhenTheActionSucceeds() {
transactionExecutor.execute(() -> {
insertDefinition("committed", "Committed definition");
return null;
});
assertEquals("Committed definition", definitionName("committed"));
}
@Test
void rollsBackAllStatementsWhenTheActionFails() {
assertThrows(IllegalStateException.class, () -> transactionExecutor.execute(() -> {
insertDefinition("rolled-back", "Rolled back definition");
throw new IllegalStateException("boom");
}));
assertNull(definitionName("rolled-back"));
}
@Test
void nestedExecutionsJoinTheSurroundingTransaction() {
assertThrows(IllegalStateException.class, () -> transactionExecutor.execute(() -> {
insertDefinition("outer", "Outer definition");
transactionExecutor.execute(() -> {
insertDefinition("inner", "Inner definition");
return null;
});
throw new IllegalStateException("boom");
}));
assertNull(definitionName("outer"));
assertNull(definitionName("inner"));
}
private void insertDefinition(String id, String name) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO ordo_process_definition (id, name) VALUES (?, ?)")) {
insert.setString(1, id);
insert.setString(2, name);
insert.executeUpdate();
} catch (SQLException e) {
throw new JdbcStorageException("failed to insert definition: " + id, e);
} finally {
connectionProvider.close(connection);
}
}
private String definitionName(String id) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(
"SELECT name FROM ordo_process_definition WHERE id = ?")) {
select.setString(1, id);
try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next() ? resultSet.getString("name") : null;
}
} catch (SQLException e) {
throw new JdbcStorageException("failed to load definition: " + id, e);
} finally {
connectionProvider.close(connection);
}
}
}