From 397a9c229b315e627a95b051a4033f6c31615803 Mon Sep 17 00:00:00 2001 From: 0264408 Date: Tue, 8 Sep 2026 16:56:04 +0800 Subject: [PATCH] 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 --- .gitignore | 39 +++ .idea/.gitignore | 10 + .idea/encodings.xml | 15 + .idea/misc.xml | 14 + .idea/vcs.xml | 7 + docs/jdbc-plan.md | 131 +++++++++ ordo-api/pom.xml | 22 ++ .../com/jetlumen/ordo/api/ApprovalStep.java | 16 ++ .../com/jetlumen/ordo/api/ApprovalTask.java | 7 + .../jetlumen/ordo/api/AssigneeResolver.java | 11 + .../com/jetlumen/ordo/api/OrdoEngine.java | 26 ++ .../com/jetlumen/ordo/api/ProcessContext.java | 20 ++ .../jetlumen/ordo/api/ProcessDefinition.java | 23 ++ .../jetlumen/ordo/api/ProcessInstance.java | 7 + .../com/jetlumen/ordo/api/ProcessStatus.java | 5 + .../com/jetlumen/ordo/api/TaskAction.java | 15 + .../com/jetlumen/ordo/api/TaskStatus.java | 5 + .../ordo/api/TransactionExecutor.java | 13 + .../DefinitionAlreadyExistsException.java | 7 + .../DefinitionNotFoundException.java | 7 + .../ordo/api/exception/OrdoException.java | 12 + .../TaskAlreadyCompletedException.java | 7 + .../api/exception/TaskNotFoundException.java | 7 + .../UnauthorizedTaskOperationException.java | 7 + .../repository/ApprovalTaskRepository.java | 27 ++ .../ProcessDefinitionRepository.java | 17 ++ .../repository/ProcessInstanceRepository.java | 16 ++ .../ordo/api/ProcessDefinitionTest.java | 50 ++++ ordo-core/pom.xml | 27 ++ .../jetlumen/ordo/core/DefaultOrdoEngine.java | 200 +++++++++++++ .../ordo/core/InMemoryOrdoEngine.java | 84 ++++++ .../ordo/core/NoopTransactionExecutor.java | 15 + .../InMemoryApprovalTaskRepository.java | 56 ++++ .../InMemoryProcessDefinitionRepository.java | 23 ++ .../InMemoryProcessInstanceRepository.java | 28 ++ .../ordo/core/InMemoryOrdoEngineTest.java | 178 ++++++++++++ ordo-example/pom.xml | 27 ++ .../ordo/example/LeaveRequestExample.java | 32 +++ ordo-storage-jdbc/pom.xml | 47 +++ .../jdbc/JdbcApprovalTaskRepository.java | 122 ++++++++ .../storage/jdbc/JdbcConnectionProvider.java | 74 +++++ .../jdbc/JdbcProcessDefinitionRepository.java | 102 +++++++ .../jdbc/JdbcProcessInstanceRepository.java | 75 +++++ .../storage/jdbc/JdbcStorageException.java | 10 + .../storage/jdbc/JdbcTransactionExecutor.java | 57 ++++ .../jdbc/mapper/ApprovalStepMapper.java | 17 ++ .../jdbc/mapper/ApprovalTaskMapper.java | 54 ++++ .../jdbc/mapper/ProcessContextCodec.java | 38 +++ .../jdbc/mapper/ProcessInstanceMapper.java | 44 +++ .../db/migration/V1__create_ordo_tables.sql | 48 ++++ .../jdbc/JdbcApprovalTaskRepositoryTest.java | 143 +++++++++ .../jdbc/JdbcOrdoEngineIntegrationTest.java | 133 +++++++++ .../jdbc/JdbcPostgresIntegrationTest.java | 271 ++++++++++++++++++ .../JdbcProcessDefinitionRepositoryTest.java | 49 ++++ .../JdbcProcessInstanceRepositoryTest.java | 74 +++++ .../ordo/storage/jdbc/JdbcTestSupport.java | 53 ++++ .../jdbc/JdbcTransactionExecutorTest.java | 88 ++++++ pom.xml | 67 +++++ 58 files changed, 2779 insertions(+) create mode 100644 .gitignore create mode 100644 .idea/.gitignore create mode 100644 .idea/encodings.xml create mode 100644 .idea/misc.xml create mode 100644 .idea/vcs.xml create mode 100644 docs/jdbc-plan.md create mode 100644 ordo-api/pom.xml create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalStep.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalTask.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/AssigneeResolver.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessContext.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessDefinition.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessInstance.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessStatus.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/TaskAction.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/TaskStatus.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/TransactionExecutor.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionAlreadyExistsException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionNotFoundException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/OrdoException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskAlreadyCompletedException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskNotFoundException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/exception/UnauthorizedTaskOperationException.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java create mode 100644 ordo-api/src/test/java/com/jetlumen/ordo/api/ProcessDefinitionTest.java create mode 100644 ordo-core/pom.xml create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/NoopTransactionExecutor.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java create mode 100644 ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java create mode 100644 ordo-example/pom.xml create mode 100644 ordo-example/src/main/java/com/jetlumen/ordo/example/LeaveRequestExample.java create mode 100644 ordo-storage-jdbc/pom.xml create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcConnectionProvider.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcStorageException.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutor.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalStepMapper.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessContextCodec.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessInstanceMapper.java create mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutorTest.java create mode 100644 pom.xml diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..480bdf5 --- /dev/null +++ b/.gitignore @@ -0,0 +1,39 @@ +target/ +!.mvn/wrapper/maven-wrapper.jar +!**/src/main/**/target/ +!**/src/test/**/target/ +.kotlin + +### IntelliJ IDEA ### +.idea/modules.xml +.idea/jarRepositories.xml +.idea/compiler.xml +.idea/libraries/ +*.iws +*.iml +*.ipr + +### Eclipse ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans +.sts4-cache + +### NetBeans ### +/nbproject/private/ +/nbbuild/ +/dist/ +/nbdist/ +/.nb-gradle/ +build/ +!**/src/main/**/build/ +!**/src/test/**/build/ + +### VS Code ### +.vscode/ + +### Mac OS ### +.DS_Store \ No newline at end of file diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..30cf57e --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,10 @@ +# Default ignored files +/shelf/ +/workspace.xml +# Editor-based HTTP Client requests +/httpRequests/ +# Ignored default folder with query files +/queries/ +# Datasource local storage ignored files +/dataSources/ +/dataSources.local.xml diff --git a/.idea/encodings.xml b/.idea/encodings.xml new file mode 100644 index 0000000..b4b2816 --- /dev/null +++ b/.idea/encodings.xml @@ -0,0 +1,15 @@ + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/.idea/misc.xml b/.idea/misc.xml new file mode 100644 index 0000000..de41ab4 --- /dev/null +++ b/.idea/misc.xml @@ -0,0 +1,14 @@ + + + + + + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..8306744 --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,7 @@ + + + + + + + \ No newline at end of file diff --git a/docs/jdbc-plan.md b/docs/jdbc-plan.md new file mode 100644 index 0000000..dd870e1 --- /dev/null +++ b/docs/jdbc-plan.md @@ -0,0 +1,131 @@ +建议将 JDBC 做成独立模块,并在实现前先补上“跨仓储事务”和“并发安全”两个能力;否则一次审批会拆成多次独立数据库操作,容易留下半完成流程。 + +1. 新建模块 + +```text +ordo-storage-jdbc/ +├── pom.xml +└── src/main/ + ├── java/com/jetlumen/ordo/storage/jdbc/ + │ ├── JdbcTransactionExecutor.java + │ ├── JdbcProcessDefinitionRepository.java + │ ├── JdbcProcessInstanceRepository.java + │ ├── JdbcApprovalTaskRepository.java + │ ├── JdbcConnectionProvider.java + │ └── mapper/ + └── resources/db/migration/ + └── V1__create_ordo_tables.sql +``` + +依赖只需要 `ordo-api`、`javax.sql.DataSource` 和 JDBC 驱动;先不要依赖 Spring。 + +2. 先补事务边界 + +在 `ordo-api` 增加一个通用端口: + +```java +public interface TransactionExecutor { + T execute(Supplier action); +} +``` + +`DefaultOrdoEngine` 的 `start`、`approve`、`reject` 应在同一个事务中执行。 + +这保证: + +- 发起流程时,“创建实例 + 创建第一条任务”要么都成功,要么都回滚。 +- 审批时,“完成旧任务 + 创建下一任务 / 结束实例”要么都成功,要么都回滚。 + +内存实现提供无操作事务执行器;JDBC 实现使用同一条线程绑定的 `Connection`,执行 `commit` 或 `rollback`。 + +3. 调整仓储 SPI 的并发语义 + +当前 `save` 是覆盖式写入,JDBC 下无法避免两个用户同时审批同一任务。建议在落 JDBC 前调整: + +```java +boolean insertIfAbsent(ProcessDefinition definition); + +boolean completeIfPending(ApprovalTask completedTask); +``` + +`completeIfPending` 对应 SQL: + +```sql +UPDATE ordo_approval_task +SET status = ?, completed_at = ?, action_actor = ?, action_comment = ?, action_at = ? +WHERE id = ? AND status = 'PENDING' +``` + +受影响行数为 `0` 时,抛出 `TaskAlreadyCompletedException`。这比仅依赖 JVM 内的 `synchronized` 更可靠。 + +4. 数据库模型 + +采用规范化表,不把步骤和任务都塞进 JSON。 + +```text +ordo_process_definition +- id varchar(64) primary key +- name varchar(255) not null + +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) + +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 null + +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 null +- action_actor varchar(255) null +- action_comment text null +- action_at timestamp null +``` + +至少建立: + +```text +ordo_approval_task(instance_id) +ordo_approval_task(status, assignee) +ordo_process_instance(definition_id) +``` + +5. `ProcessContext` 的持久化 + +`ProcessContext.variables` 适合存为 `context_json`。在 JDBC 模块内部使用 Jackson 做序列化与反序列化,不要让 `ordo-api` 依赖 Jackson。 + +v0.1 可以约定上下文仅支持 JSON 兼容值:字符串、数字、布尔值、列表、嵌套 Map。日期、枚举和自定义 Java 对象以后再通过可插拔 `ContextCodec` 解决。 + +6. JDBC 实现顺序 + +- 先写 `V1__create_ordo_tables.sql` +- 实现 `JdbcTransactionExecutor` +- 实现定义仓储与步骤读写 +- 实现实例仓储及 `context_json` +- 实现任务仓储及待办查询 +- 调整 `DefaultOrdoEngine` 使用事务和条件更新 +- 为 JDBC 仓储添加集成测试 + +7. 测试策略 + +先用 H2 快速验证 CRUD 和映射;并发条件更新、时间类型、唯一约束等最终应使用 Testcontainers 加 PostgreSQL 或 MySQL 验证。 + +建议第一版目标是 PostgreSQL;表结构、`timestamp` 语义和 JSON 支持都会更明确。等 JDBC 实现稳定后,再创建 Spring Boot Starter:Starter 只负责注入 `DataSource`、JDBC 仓储、事务执行器和 `DefaultOrdoEngine`。 + +帮忙按照这个实现一下JDBC模块 diff --git a/ordo-api/pom.xml b/ordo-api/pom.xml new file mode 100644 index 0000000..34afa63 --- /dev/null +++ b/ordo-api/pom.xml @@ -0,0 +1,22 @@ + + + 4.0.0 + + com.jetlumen + ordo + 1.0-SNAPSHOT + + + ordo-api + + + + org.junit.jupiter + junit-jupiter + test + + + + diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalStep.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalStep.java new file mode 100644 index 0000000..25df328 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalStep.java @@ -0,0 +1,16 @@ +package com.jetlumen.ordo.api; + +/** A single, named approval step in a linear process definition. */ +public record ApprovalStep(String id, String name, String assignee) { + public ApprovalStep { + requireText(id, "step id"); + requireText(name, "step name"); + requireText(assignee, "step assignee"); + } + + static void requireText(String value, String field) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(field + " must not be blank"); + } + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalTask.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalTask.java new file mode 100644 index 0000000..8fe8b56 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ApprovalTask.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api; + +import java.time.Instant; + +public record ApprovalTask(String id, String instanceId, String stepId, String name, String assignee, + TaskStatus status, Instant createdAt, Instant completedAt, TaskAction action) { +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/AssigneeResolver.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/AssigneeResolver.java new file mode 100644 index 0000000..e549021 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/AssigneeResolver.java @@ -0,0 +1,11 @@ +package com.jetlumen.ordo.api; + +/** Resolves the current assignee for an approval step when a task is created. */ +@FunctionalInterface +public interface AssigneeResolver { + String resolve(ApprovalStep step, ProcessContext context); + + static AssigneeResolver direct() { + return (step, context) -> step.assignee(); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java new file mode 100644 index 0000000..ba19f56 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java @@ -0,0 +1,26 @@ +package com.jetlumen.ordo.api; + +import java.util.List; +import java.util.Optional; + +/** Public entry point for definition registration and approval operations. */ +public interface OrdoEngine { + void register(ProcessDefinition definition); + default ProcessInstance start(String definitionId, String initiator) { + return start(definitionId, initiator, ProcessContext.empty()); + } + ProcessInstance start(String definitionId, String initiator, ProcessContext context); + default ApprovalTask approve(String taskId, String actor) { + return approve(taskId, actor, null); + } + ApprovalTask approve(String taskId, String actor, String comment); + default ApprovalTask reject(String taskId, String actor) { + return reject(taskId, actor, null); + } + ApprovalTask reject(String taskId, String actor, String comment); + Optional findInstance(String instanceId); + Optional findTask(String taskId); + List findTasks(String instanceId); + List findPendingTasksByAssignee(String assignee); + List findPendingTasksByInstanceId(String instanceId); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessContext.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessContext.java new file mode 100644 index 0000000..789dd47 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessContext.java @@ -0,0 +1,20 @@ +package com.jetlumen.ordo.api; + +import java.util.Map; +import java.util.Objects; +import java.util.Optional; + +/** Immutable business data available while a process instance is running. */ +public record ProcessContext(Map variables) { + public ProcessContext { + variables = Map.copyOf(Objects.requireNonNull(variables, "variables must not be null")); + } + + public static ProcessContext empty() { + return new ProcessContext(Map.of()); + } + + public Optional value(String name) { + return Optional.ofNullable(variables.get(name)); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessDefinition.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessDefinition.java new file mode 100644 index 0000000..24fcfc4 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessDefinition.java @@ -0,0 +1,23 @@ +package com.jetlumen.ordo.api; + +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +/** Immutable blueprint for a linear approval process. */ +public record ProcessDefinition(String id, String name, List steps) { + public ProcessDefinition { + ApprovalStep.requireText(id, "definition id"); + ApprovalStep.requireText(name, "definition name"); + steps = List.copyOf(steps); + if (steps.isEmpty()) { + throw new IllegalArgumentException("a definition must contain at least one approval step"); + } + Set ids = new HashSet<>(); + for (ApprovalStep step : steps) { + if (!ids.add(step.id())) { + throw new IllegalArgumentException("duplicate step id: " + step.id()); + } + } + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessInstance.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessInstance.java new file mode 100644 index 0000000..bab7381 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessInstance.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api; + +import java.time.Instant; + +public record ProcessInstance(String id, String definitionId, String initiator, ProcessStatus status, + Instant startedAt, Instant finishedAt, ProcessContext context) { +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessStatus.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessStatus.java new file mode 100644 index 0000000..db61a15 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessStatus.java @@ -0,0 +1,5 @@ +package com.jetlumen.ordo.api; + +public enum ProcessStatus { + RUNNING, APPROVED, REJECTED +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskAction.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskAction.java new file mode 100644 index 0000000..f5d12c9 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskAction.java @@ -0,0 +1,15 @@ +package com.jetlumen.ordo.api; + +import java.time.Instant; +import java.util.Objects; + +/** Immutable audit record created when an approval task is completed. */ +public record TaskAction(String actor, String comment, Instant operatedAt) { + public TaskAction { + if (actor == null || actor.isBlank()) { + throw new IllegalArgumentException("actor must not be blank"); + } + comment = comment == null || comment.isBlank() ? null : comment.strip(); + Objects.requireNonNull(operatedAt, "operatedAt must not be null"); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskStatus.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskStatus.java new file mode 100644 index 0000000..9535d08 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/TaskStatus.java @@ -0,0 +1,5 @@ +package com.jetlumen.ordo.api; + +public enum TaskStatus { + PENDING, APPROVED, REJECTED +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/TransactionExecutor.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/TransactionExecutor.java new file mode 100644 index 0000000..021d4b7 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/TransactionExecutor.java @@ -0,0 +1,13 @@ +package com.jetlumen.ordo.api; + +import java.util.function.Supplier; + +/** + * Port for the transaction boundary around multistep write operations. + * Implementations either run the action as-is (in-memory storage) or commit + * and roll back atomically (JDBC storage). + */ +@FunctionalInterface +public interface TransactionExecutor { + T execute(Supplier action); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionAlreadyExistsException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionAlreadyExistsException.java new file mode 100644 index 0000000..5165f1d --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionAlreadyExistsException.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api.exception; + +public final class DefinitionAlreadyExistsException extends OrdoException { + public DefinitionAlreadyExistsException(String definitionId) { + super("definition already exists: " + definitionId); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionNotFoundException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionNotFoundException.java new file mode 100644 index 0000000..6504947 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/DefinitionNotFoundException.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api.exception; + +public final class DefinitionNotFoundException extends OrdoException { + public DefinitionNotFoundException(String definitionId) { + super("definition not found: " + definitionId); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/OrdoException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/OrdoException.java new file mode 100644 index 0000000..994eba8 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/OrdoException.java @@ -0,0 +1,12 @@ +package com.jetlumen.ordo.api.exception; + +/** Base type for business-rule violations raised by the Ordo runtime. */ +public class OrdoException extends RuntimeException { + public OrdoException(String message) { + super(message); + } + + public OrdoException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskAlreadyCompletedException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskAlreadyCompletedException.java new file mode 100644 index 0000000..b603e63 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskAlreadyCompletedException.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api.exception; + +public final class TaskAlreadyCompletedException extends OrdoException { + public TaskAlreadyCompletedException(String taskId) { + super("task is already completed: " + taskId); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskNotFoundException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskNotFoundException.java new file mode 100644 index 0000000..a16e1a0 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/TaskNotFoundException.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api.exception; + +public final class TaskNotFoundException extends OrdoException { + public TaskNotFoundException(String taskId) { + super("task not found: " + taskId); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/UnauthorizedTaskOperationException.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/UnauthorizedTaskOperationException.java new file mode 100644 index 0000000..f31f08f --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/exception/UnauthorizedTaskOperationException.java @@ -0,0 +1,7 @@ +package com.jetlumen.ordo.api.exception; + +public final class UnauthorizedTaskOperationException extends OrdoException { + public UnauthorizedTaskOperationException(String taskId, String actor) { + super("actor '" + actor + "' is not the assignee for task: " + taskId); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java new file mode 100644 index 0000000..fab3f04 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java @@ -0,0 +1,27 @@ +package com.jetlumen.ordo.api.repository; + +import com.jetlumen.ordo.api.ApprovalTask; + +import java.util.List; +import java.util.Optional; + +/** Storage port for approval tasks and their pending-task indexes. */ +public interface ApprovalTaskRepository { + /** Inserts a new task. */ + void save(ApprovalTask task); + + Optional findById(String taskId); + + List findByInstanceId(String instanceId); + + List findPendingByAssignee(String assignee); + + List findPendingByInstanceId(String instanceId); + + /** + * Atomically completes the task only if it is still pending. + * + * @return true if the update was applied, false if the task had already been completed by a concurrent operation + */ + boolean completeIfPending(ApprovalTask completedTask); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java new file mode 100644 index 0000000..7e1ec4b --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java @@ -0,0 +1,17 @@ +package com.jetlumen.ordo.api.repository; + +import com.jetlumen.ordo.api.ProcessDefinition; + +import java.util.Optional; + +/** Storage port for immutable process definitions. */ +public interface ProcessDefinitionRepository { + /** + * Inserts the definition if no definition with the same id exists. + * + * @return true if the definition was inserted, false if a definition with the same id already exists + */ + boolean insertIfAbsent(ProcessDefinition definition); + + Optional findById(String definitionId); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java new file mode 100644 index 0000000..7f2db0a --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java @@ -0,0 +1,16 @@ +package com.jetlumen.ordo.api.repository; + +import com.jetlumen.ordo.api.ProcessInstance; + +import java.util.Optional; + +/** Storage port for process instances. */ +public interface ProcessInstanceRepository { + /** Inserts a new instance; fails if the id already exists. */ + void insert(ProcessInstance instance); + + /** Replaces the stored instance with the same id, e.g. when the process reaches a terminal status. */ + void update(ProcessInstance instance); + + Optional findById(String instanceId); +} diff --git a/ordo-api/src/test/java/com/jetlumen/ordo/api/ProcessDefinitionTest.java b/ordo-api/src/test/java/com/jetlumen/ordo/api/ProcessDefinitionTest.java new file mode 100644 index 0000000..3f28fcf --- /dev/null +++ b/ordo-api/src/test/java/com/jetlumen/ordo/api/ProcessDefinitionTest.java @@ -0,0 +1,50 @@ +package com.jetlumen.ordo.api; + +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class ProcessDefinitionTest { + + @Test + void rejectsBlankStepFields() { + assertThrows(IllegalArgumentException.class, () -> new ApprovalStep(" ", "Manager", "maria")); + assertThrows(IllegalArgumentException.class, () -> new ApprovalStep("manager", " ", "maria")); + assertThrows(IllegalArgumentException.class, () -> new ApprovalStep("manager", "Manager", " ")); + } + + @Test + void rejectsDefinitionsWithoutStepsOrWithDuplicateStepIds() { + assertThrows(IllegalArgumentException.class, + () -> new ProcessDefinition("leave", "Leave request", List.of())); + assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition("leave", "Leave request", List.of( + new ApprovalStep("manager", "Manager approval", "maria"), + new ApprovalStep("manager", "HR approval", "henry") + ))); + } + + @Test + void rejectsBlankDefinitionFields() { + List steps = List.of(new ApprovalStep("manager", "Manager approval", "maria")); + + assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition(" ", "Leave request", steps)); + assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition("leave", " ", steps)); + } + + @Test + void copiesTheSuppliedStepList() { + List suppliedSteps = new ArrayList<>(); + suppliedSteps.add(new ApprovalStep("manager", "Manager approval", "maria")); + ProcessDefinition definition = new ProcessDefinition("leave", "Leave request", suppliedSteps); + + suppliedSteps.add(new ApprovalStep("hr", "HR approval", "henry")); + + assertEquals(1, definition.steps().size()); + assertThrows(UnsupportedOperationException.class, + () -> definition.steps().add(new ApprovalStep("lead", "Lead approval", "lee"))); + } +} diff --git a/ordo-core/pom.xml b/ordo-core/pom.xml new file mode 100644 index 0000000..faedace --- /dev/null +++ b/ordo-core/pom.xml @@ -0,0 +1,27 @@ + + + 4.0.0 + + com.jetlumen + ordo + 1.0-SNAPSHOT + + + ordo-core + + + + com.jetlumen + ordo-api + ${project.version} + + + org.junit.jupiter + junit-jupiter + test + + + + diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java new file mode 100644 index 0000000..5a63ce4 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java @@ -0,0 +1,200 @@ +package com.jetlumen.ordo.core; + +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.api.TransactionExecutor; +import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException; +import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; +import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException; +import com.jetlumen.ordo.api.exception.TaskNotFoundException; +import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; +import com.jetlumen.ordo.api.repository.ApprovalTaskRepository; +import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; +import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; + +import java.time.Clock; +import java.time.Instant; +import java.util.List; +import java.util.Objects; +import java.util.Optional; +import java.util.UUID; + +/** + * Repository-backed implementation of the v0.1 linear approval runtime. + * Multistep write operations run inside a {@link TransactionExecutor} and + * rely on conditional repository updates, so the engine stays correct even + * when several JVMs share the same storage. + */ +public final class DefaultOrdoEngine implements OrdoEngine { + private final Clock clock; + private final AssigneeResolver assigneeResolver; + private final TransactionExecutor transactionExecutor; + private final ProcessDefinitionRepository definitionRepository; + private final ProcessInstanceRepository instanceRepository; + private final ApprovalTaskRepository taskRepository; + + public DefaultOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, TransactionExecutor transactionExecutor, + ProcessDefinitionRepository definitionRepository, + ProcessInstanceRepository instanceRepository, + ApprovalTaskRepository taskRepository) { + this.clock = Objects.requireNonNull(clock, "clock must not be null"); + this.assigneeResolver = Objects.requireNonNull(assigneeResolver, "assigneeResolver must not be null"); + this.transactionExecutor = Objects.requireNonNull(transactionExecutor, "transactionExecutor must not be null"); + this.definitionRepository = Objects.requireNonNull(definitionRepository, "definitionRepository must not be null"); + this.instanceRepository = Objects.requireNonNull(instanceRepository, "instanceRepository must not be null"); + this.taskRepository = Objects.requireNonNull(taskRepository, "taskRepository must not be null"); + } + + @Override + public synchronized void register(ProcessDefinition definition) { + Objects.requireNonNull(definition, "definition must not be null"); + transactionExecutor.execute(() -> { + if (!definitionRepository.insertIfAbsent(definition)) { + throw new DefinitionAlreadyExistsException(definition.id()); + } + return null; + }); + } + + @Override + public synchronized ProcessInstance start(String definitionId, String initiator, ProcessContext context) { + requireText(initiator, "initiator"); + Objects.requireNonNull(context, "context must not be null"); + return transactionExecutor.execute(() -> { + ProcessDefinition definition = requireDefinition(definitionId); + Instant now = clock.instant(); + ProcessInstance instance = new ProcessInstance(nextId(), definition.id(), initiator, + ProcessStatus.RUNNING, now, null, context); + instanceRepository.insert(instance); + createTask(instance, definition.steps().getFirst(), now); + return instance; + }); + } + + @Override + public synchronized ApprovalTask approve(String taskId, String actor, String comment) { + return transactionExecutor.execute(() -> { + ApprovalTask task = requirePendingTaskForActor(taskId, actor); + Instant now = clock.instant(); + ApprovalTask completedTask = completeTask(task, TaskStatus.APPROVED, new TaskAction(actor, comment, now)); + ProcessInstance instance = requireInstance(task.instanceId()); + ProcessDefinition definition = requireDefinition(instance.definitionId()); + int stepIndex = indexOf(definition, task.stepId()); + if (stepIndex == definition.steps().size() - 1) { + completeInstance(instance, ProcessStatus.APPROVED, now); + } else { + createTask(instance, definition.steps().get(stepIndex + 1), now); + } + return completedTask; + }); + } + + @Override + public synchronized ApprovalTask reject(String taskId, String actor, String comment) { + return transactionExecutor.execute(() -> { + ApprovalTask task = requirePendingTaskForActor(taskId, actor); + Instant now = clock.instant(); + ApprovalTask completedTask = completeTask(task, TaskStatus.REJECTED, new TaskAction(actor, comment, now)); + completeInstance(requireInstance(task.instanceId()), ProcessStatus.REJECTED, now); + return completedTask; + }); + } + + @Override + public synchronized Optional findInstance(String instanceId) { + return instanceRepository.findById(instanceId); + } + + @Override + public synchronized Optional findTask(String taskId) { + return taskRepository.findById(taskId); + } + + @Override + public synchronized List findTasks(String instanceId) { + return taskRepository.findByInstanceId(instanceId); + } + + @Override + public synchronized List findPendingTasksByAssignee(String assignee) { + requireText(assignee, "assignee"); + return taskRepository.findPendingByAssignee(assignee); + } + + @Override + public synchronized List findPendingTasksByInstanceId(String instanceId) { + requireText(instanceId, "instance id"); + return taskRepository.findPendingByInstanceId(instanceId); + } + + private void createTask(ProcessInstance instance, ApprovalStep step, Instant now) { + String assignee = assigneeResolver.resolve(step, instance.context()); + requireText(assignee, "resolved assignee"); + taskRepository.save(new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee, + TaskStatus.PENDING, now, null, null)); + } + + private ApprovalTask requirePendingTaskForActor(String taskId, String actor) { + requireText(actor, "actor"); + ApprovalTask task = taskRepository.findById(taskId) + .orElseThrow(() -> new TaskNotFoundException(taskId)); + if (task.status() != TaskStatus.PENDING) { + throw new TaskAlreadyCompletedException(taskId); + } + if (!task.assignee().equals(actor)) { + throw new UnauthorizedTaskOperationException(taskId, actor); + } + return task; + } + + private ApprovalTask completeTask(ApprovalTask task, TaskStatus status, TaskAction action) { + ApprovalTask completed = new ApprovalTask(task.id(), task.instanceId(), task.stepId(), task.name(), + task.assignee(), status, task.createdAt(), action.operatedAt(), action); + if (!taskRepository.completeIfPending(completed)) { + throw new TaskAlreadyCompletedException(task.id()); + } + return completed; + } + + private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now) { + instanceRepository.update(new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(), + status, instance.startedAt(), now, instance.context())); + } + + private ProcessDefinition requireDefinition(String definitionId) { + return definitionRepository.findById(definitionId) + .orElseThrow(() -> new DefinitionNotFoundException(definitionId)); + } + + private ProcessInstance requireInstance(String instanceId) { + return instanceRepository.findById(instanceId) + .orElseThrow(() -> new IllegalStateException("instance not found: " + instanceId)); + } + + private static int indexOf(ProcessDefinition definition, String stepId) { + for (int index = 0; index < definition.steps().size(); index++) { + if (definition.steps().get(index).id().equals(stepId)) { + return index; + } + } + throw new IllegalStateException("step not found in definition: " + stepId); + } + + private static String nextId() { + return UUID.randomUUID().toString(); + } + + private static void requireText(String value, String name) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(name + " must not be blank"); + } + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java new file mode 100644 index 0000000..310ac89 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java @@ -0,0 +1,84 @@ +package com.jetlumen.ordo.core; + +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.core.repository.InMemoryApprovalTaskRepository; +import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository; +import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository; + +import java.time.Clock; +import java.util.List; +import java.util.Optional; + +/** Development and test engine backed by the in-memory repository implementations. */ +public final class InMemoryOrdoEngine implements OrdoEngine { + private final DefaultOrdoEngine delegate; + + public InMemoryOrdoEngine() { + this(Clock.systemUTC(), AssigneeResolver.direct()); + } + + public InMemoryOrdoEngine(Clock clock) { + this(clock, AssigneeResolver.direct()); + } + + public InMemoryOrdoEngine(AssigneeResolver assigneeResolver) { + this(Clock.systemUTC(), assigneeResolver); + } + + public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver) { + this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, new NoopTransactionExecutor(), + new InMemoryProcessDefinitionRepository(), + new InMemoryProcessInstanceRepository(), + new InMemoryApprovalTaskRepository()); + } + + @Override + public void register(ProcessDefinition definition) { + delegate.register(definition); + } + + @Override + public ProcessInstance start(String definitionId, String initiator, ProcessContext context) { + return delegate.start(definitionId, initiator, context); + } + + @Override + public ApprovalTask approve(String taskId, String actor, String comment) { + return delegate.approve(taskId, actor, comment); + } + + @Override + public ApprovalTask reject(String taskId, String actor, String comment) { + return delegate.reject(taskId, actor, comment); + } + + @Override + public Optional findInstance(String instanceId) { + return delegate.findInstance(instanceId); + } + + @Override + public Optional findTask(String taskId) { + return delegate.findTask(taskId); + } + + @Override + public List findTasks(String instanceId) { + return delegate.findTasks(instanceId); + } + + @Override + public List findPendingTasksByAssignee(String assignee) { + return delegate.findPendingTasksByAssignee(assignee); + } + + @Override + public List findPendingTasksByInstanceId(String instanceId) { + return delegate.findPendingTasksByInstanceId(instanceId); + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/NoopTransactionExecutor.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/NoopTransactionExecutor.java new file mode 100644 index 0000000..eafc7f9 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/NoopTransactionExecutor.java @@ -0,0 +1,15 @@ +package com.jetlumen.ordo.core; + +import com.jetlumen.ordo.api.TransactionExecutor; + +import java.util.Objects; +import java.util.function.Supplier; + +/** Transaction boundary that runs the action as-is; used by the in-memory engine. */ +public final class NoopTransactionExecutor implements TransactionExecutor { + @Override + public T execute(Supplier action) { + Objects.requireNonNull(action, "action must not be null"); + return action.get(); + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java new file mode 100644 index 0000000..612f0c4 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java @@ -0,0 +1,56 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ApprovalTask; +import com.jetlumen.ordo.api.TaskStatus; +import com.jetlumen.ordo.api.repository.ApprovalTaskRepository; + +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; + +/** Development-only in-memory implementation of the task storage port. */ +public final class InMemoryApprovalTaskRepository implements ApprovalTaskRepository { + private final Map tasks = new LinkedHashMap<>(); + + @Override + public synchronized void save(ApprovalTask task) { + tasks.put(task.id(), task); + } + + @Override + public synchronized Optional findById(String taskId) { + return Optional.ofNullable(tasks.get(taskId)); + } + + @Override + public synchronized List findByInstanceId(String instanceId) { + return tasks.values().stream().filter(task -> task.instanceId().equals(instanceId)).toList(); + } + + @Override + public synchronized List findPendingByAssignee(String assignee) { + return tasks.values().stream() + .filter(task -> task.status() == TaskStatus.PENDING) + .filter(task -> task.assignee().equals(assignee)) + .toList(); + } + + @Override + public synchronized List findPendingByInstanceId(String instanceId) { + return tasks.values().stream() + .filter(task -> task.status() == TaskStatus.PENDING) + .filter(task -> task.instanceId().equals(instanceId)) + .toList(); + } + + @Override + public synchronized boolean completeIfPending(ApprovalTask completedTask) { + ApprovalTask current = tasks.get(completedTask.id()); + if (current == null || current.status() != TaskStatus.PENDING) { + return false; + } + tasks.put(completedTask.id(), completedTask); + return true; + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java new file mode 100644 index 0000000..25e71f9 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java @@ -0,0 +1,23 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; + +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; + +/** Development-only in-memory implementation of the definition storage port. */ +public final class InMemoryProcessDefinitionRepository implements ProcessDefinitionRepository { + private final Map definitions = new HashMap<>(); + + @Override + public synchronized boolean insertIfAbsent(ProcessDefinition definition) { + return definitions.putIfAbsent(definition.id(), definition) == null; + } + + @Override + public synchronized Optional findById(String definitionId) { + return Optional.ofNullable(definitions.get(definitionId)); + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java new file mode 100644 index 0000000..a746d2e --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java @@ -0,0 +1,28 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; + +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; + +/** Development-only in-memory implementation of the instance storage port. */ +public final class InMemoryProcessInstanceRepository implements ProcessInstanceRepository { + private final Map instances = new HashMap<>(); + + @Override + public synchronized void insert(ProcessInstance instance) { + instances.put(instance.id(), instance); + } + + @Override + public synchronized void update(ProcessInstance instance) { + instances.put(instance.id(), instance); + } + + @Override + public synchronized Optional findById(String instanceId) { + return Optional.ofNullable(instances.get(instanceId)); + } +} diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java new file mode 100644 index 0000000..22788af --- /dev/null +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java @@ -0,0 +1,178 @@ +package com.jetlumen.ordo.core; + +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.ProcessStatus; +import com.jetlumen.ordo.api.TaskStatus; +import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException; +import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; +import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException; +import com.jetlumen.ordo.api.exception.TaskNotFoundException; +import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class InMemoryOrdoEngineTest { + private InMemoryOrdoEngine engine; + + @BeforeEach + void setUp() { + engine = new InMemoryOrdoEngine(); + engine.register(new ProcessDefinition("leave", "Leave request", List.of( + new ApprovalStep("manager", "Manager approval", "maria"), + new ApprovalStep("hr", "HR approval", "henry") + ))); + } + + @Test + void completesASequentialApprovalProcess() { + var instance = engine.start("leave", "alice"); + ApprovalTask managerTask = engine.findTasks(instance.id()).getFirst(); + + assertEquals(TaskStatus.APPROVED, engine.approve(managerTask.id(), "maria").status()); + ApprovalTask hrTask = engine.findTasks(instance.id()).get(1); + assertEquals("henry", hrTask.assignee()); + engine.approve(hrTask.id(), "henry"); + + assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status()); + } + + @Test + void rejectionTerminatesTheProcess() { + var instance = engine.start("leave", "alice"); + ApprovalTask task = engine.findTasks(instance.id()).getFirst(); + + ApprovalTask rejectedTask = engine.reject(task.id(), "maria", "Insufficient leave balance"); + + assertEquals(TaskStatus.REJECTED, rejectedTask.status()); + assertEquals("maria", rejectedTask.action().actor()); + assertEquals("Insufficient leave balance", rejectedTask.action().comment()); + assertEquals(rejectedTask.completedAt(), rejectedTask.action().operatedAt()); + 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 rejectsDuplicateAndUnauthorizedOperations() { + var instance = engine.start("leave", "alice"); + ApprovalTask task = engine.findTasks(instance.id()).getFirst(); + + assertThrows(UnauthorizedTaskOperationException.class, () -> engine.approve(task.id(), "mallory")); + engine.approve(task.id(), "maria"); + assertThrows(TaskAlreadyCompletedException.class, () -> engine.approve(task.id(), "maria")); + } + + @Test + void exposesSpecificExceptionsForMissingAndDuplicateResources() { + assertThrows(DefinitionNotFoundException.class, () -> engine.start("missing", "alice")); + assertThrows(TaskNotFoundException.class, () -> engine.approve("missing", "maria")); + assertThrows(DefinitionAlreadyExistsException.class, () -> engine.register(new ProcessDefinition( + "leave", "Another leave request", List.of(new ApprovalStep("lead", "Lead approval", "lee"))))); + } + + @Test + void rejectsBlankRuntimeArguments() { + assertThrows(IllegalArgumentException.class, () -> engine.start("leave", " ")); + assertThrows(DefinitionNotFoundException.class, () -> engine.start(" ", "alice")); + + var instance = engine.start("leave", "alice"); + ApprovalTask task = engine.findTasks(instance.id()).getFirst(); + assertThrows(IllegalArgumentException.class, () -> engine.approve(task.id(), " ")); + assertThrows(IllegalArgumentException.class, () -> engine.reject(task.id(), " ")); + } + + @Test + void keepsSeparateInstancesIndependent() { + var aliceInstance = engine.start("leave", "alice"); + var bobInstance = engine.start("leave", "bob"); + ApprovalTask aliceTask = engine.findTasks(aliceInstance.id()).getFirst(); + + engine.approve(aliceTask.id(), "maria"); + + assertEquals(2, engine.findTasks(aliceInstance.id()).size()); + assertEquals(1, engine.findTasks(bobInstance.id()).size()); + assertEquals(ProcessStatus.RUNNING, engine.findInstance(aliceInstance.id()).orElseThrow().status()); + assertEquals(ProcessStatus.RUNNING, engine.findInstance(bobInstance.id()).orElseThrow().status()); + assertEquals("maria", engine.findTasks(bobInstance.id()).getFirst().assignee()); + assertTrue(engine.findTasks("missing-instance").isEmpty()); + } + + @Test + void findsOnlyCurrentPendingTasksByAssigneeAndInstance() { + var aliceInstance = engine.start("leave", "alice"); + var bobInstance = engine.start("leave", "bob"); + ApprovalTask aliceManagerTask = engine.findPendingTasksByInstanceId(aliceInstance.id()).getFirst(); + + assertEquals(2, engine.findPendingTasksByAssignee("maria").size()); + assertEquals(List.of(aliceManagerTask), engine.findPendingTasksByInstanceId(aliceInstance.id())); + assertTrue(engine.findPendingTasksByAssignee("henry").isEmpty()); + + engine.approve(aliceManagerTask.id(), "maria"); + ApprovalTask hrTask = engine.findPendingTasksByInstanceId(aliceInstance.id()).getFirst(); + + assertEquals("hr", hrTask.stepId()); + assertEquals("henry", hrTask.assignee()); + assertEquals(1, engine.findPendingTasksByAssignee("maria").size()); + assertEquals(List.of(hrTask), engine.findPendingTasksByAssignee("henry")); + + engine.reject(hrTask.id(), "henry"); + assertTrue(engine.findPendingTasksByInstanceId(aliceInstance.id()).isEmpty()); + assertTrue(engine.findPendingTasksByAssignee("henry").isEmpty()); + assertEquals(1, engine.findPendingTasksByInstanceId(bobInstance.id()).size()); + } + + @Test + void recordsApprovalAuditDataAndLeavesPendingTasksWithoutAnAction() { + var instance = engine.start("leave", "alice"); + ApprovalTask managerTask = engine.findPendingTasksByInstanceId(instance.id()).getFirst(); + + assertNull(managerTask.action()); + ApprovalTask approvedTask = engine.approve(managerTask.id(), "maria", " Approved for the requested dates. "); + + assertEquals("maria", approvedTask.action().actor()); + assertEquals("Approved for the requested dates.", approvedTask.action().comment()); + assertEquals(approvedTask.completedAt(), approvedTask.action().operatedAt()); + assertNotNull(approvedTask.action().operatedAt()); + } + + @Test + void rejectsBlankPendingTaskQueryArguments() { + assertThrows(IllegalArgumentException.class, () -> engine.findPendingTasksByAssignee(" ")); + assertThrows(IllegalArgumentException.class, () -> engine.findPendingTasksByInstanceId(" ")); + } + + @Test + void resolvesAssigneesFromTheProcessContext() { + InMemoryOrdoEngine contextAwareEngine = new InMemoryOrdoEngine((step, context) -> context.value(step.id()) + .filter(String.class::isInstance) + .map(String.class::cast) + .orElse(step.assignee())); + contextAwareEngine.register(new ProcessDefinition("leave", "Leave request", List.of( + new ApprovalStep("manager", "Manager approval", "maria"), + new ApprovalStep("hr", "HR approval", "henry") + ))); + + var instance = contextAwareEngine.start("leave", "alice", new ProcessContext(Map.of( + "manager", "david", + "hr", "helena" + ))); + ApprovalTask managerTask = contextAwareEngine.findPendingTasksByInstanceId(instance.id()).getFirst(); + + assertEquals("david", managerTask.assignee()); + contextAwareEngine.approve(managerTask.id(), "david"); + assertEquals("helena", contextAwareEngine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee()); + assertEquals("david", instance.context().value("manager").orElseThrow()); + } +} diff --git a/ordo-example/pom.xml b/ordo-example/pom.xml new file mode 100644 index 0000000..fa03581 --- /dev/null +++ b/ordo-example/pom.xml @@ -0,0 +1,27 @@ + + + 4.0.0 + + com.jetlumen + ordo + 1.0-SNAPSHOT + + + ordo-example + + + + com.jetlumen + ordo-api + ${project.version} + + + com.jetlumen + ordo-core + ${project.version} + + + + diff --git a/ordo-example/src/main/java/com/jetlumen/ordo/example/LeaveRequestExample.java b/ordo-example/src/main/java/com/jetlumen/ordo/example/LeaveRequestExample.java new file mode 100644 index 0000000..2e3a14a --- /dev/null +++ b/ordo-example/src/main/java/com/jetlumen/ordo/example/LeaveRequestExample.java @@ -0,0 +1,32 @@ +package com.jetlumen.ordo.example; + +import com.jetlumen.ordo.api.ApprovalStep; +import com.jetlumen.ordo.api.OrdoEngine; +import com.jetlumen.ordo.api.ProcessContext; +import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.core.InMemoryOrdoEngine; + +import java.util.List; +import java.util.Map; + +public final class LeaveRequestExample { + private LeaveRequestExample() { + } + + public static void main(String[] args) { + OrdoEngine ordo = new InMemoryOrdoEngine(); + ordo.register(new ProcessDefinition("leave-request", "Leave request", List.of( + new ApprovalStep("manager", "Manager approval", "maria"), + new ApprovalStep("hr", "HR approval", "henry") + ))); + + var instance = ordo.start("leave-request", "alice", new ProcessContext(Map.of( + "requestId", "LEAVE-2026-001" + ))); + var managerTask = ordo.findTasks(instance.id()).getFirst(); + ordo.approve(managerTask.id(), "maria", "Approved by the manager."); + var hrTask = ordo.findTasks(instance.id()).get(1); + ordo.approve(hrTask.id(), "henry", "Leave record updated."); + + } +} diff --git a/ordo-storage-jdbc/pom.xml b/ordo-storage-jdbc/pom.xml new file mode 100644 index 0000000..e2a9b06 --- /dev/null +++ b/ordo-storage-jdbc/pom.xml @@ -0,0 +1,47 @@ + + + 4.0.0 + + com.jetlumen + ordo + 1.0-SNAPSHOT + + + ordo-storage-jdbc + + + + com.jetlumen + ordo-api + ${project.version} + + + com.jetlumen + ordo-core + ${project.version} + test + + + com.fasterxml.jackson.core + jackson-databind + + + org.postgresql + postgresql + runtime + + + org.junit.jupiter + junit-jupiter + test + + + com.h2database + h2 + test + + + + diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java new file mode 100644 index 0000000..4f6919a --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java @@ -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 findById(String taskId) { + return findOne(SELECT_TASK, taskId); + } + + @Override + public List findByInstanceId(String instanceId) { + return findAll(SELECT_BY_INSTANCE, instanceId); + } + + @Override + public List findPendingByAssignee(String assignee) { + return findAll(SELECT_PENDING_BY_ASSIGNEE, assignee); + } + + @Override + public List 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 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 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 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); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcConnectionProvider.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcConnectionProvider.java new file mode 100644 index 0000000..24a18b9 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcConnectionProvider.java @@ -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 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); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java new file mode 100644 index 0000000..34384d2 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java @@ -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 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 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()); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java new file mode 100644 index 0000000..91c3ddf --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java @@ -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 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); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcStorageException.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcStorageException.java new file mode 100644 index 0000000..b041015 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcStorageException.java @@ -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); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutor.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutor.java new file mode 100644 index 0000000..1223f8c --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutor.java @@ -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 execute(Supplier 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); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalStepMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalStepMapper.java new file mode 100644 index 0000000..4f557ff --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalStepMapper.java @@ -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")); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java new file mode 100644 index 0000000..873bb44 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java @@ -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); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessContextCodec.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessContextCodec.java new file mode 100644 index 0000000..000c602 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessContextCodec.java @@ -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 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); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessInstanceMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessInstanceMapper.java new file mode 100644 index 0000000..6bced8f --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessInstanceMapper.java @@ -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"))); + } +} diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql new file mode 100644 index 0000000..bd9d011 --- /dev/null +++ b/ordo-storage-jdbc/src/main/resources/db/migration/V1__create_ordo_tables.sql @@ -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); diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java new file mode 100644 index 0000000..3eb690d --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java @@ -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)); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java new file mode 100644 index 0000000..281ffb5 --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java @@ -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)); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java new file mode 100644 index 0000000..7d5da5c --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java @@ -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. + * + *

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)); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java new file mode 100644 index 0000000..5c07ebc --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java @@ -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"))); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java new file mode 100644 index 0000000..5982574 --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java @@ -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()))); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java new file mode 100644 index 0000000..21a48ad --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java @@ -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); + } + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutorTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutorTest.java new file mode 100644 index 0000000..04da5bf --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTransactionExecutorTest.java @@ -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); + } + } +} diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..87dbdba --- /dev/null +++ b/pom.xml @@ -0,0 +1,67 @@ + + + 4.0.0 + + com.jetlumen + ordo + 1.0-SNAPSHOT + + pom + + ordo + Lightweight approval workflow engine + + ordo-core + ordo-api + ordo-example + ordo-storage-jdbc + + + + 21 + 21 + UTF-8 + 5.11.4 + 2.18.2 + 2.3.232 + 42.7.5 + + + + + + org.junit.jupiter + junit-jupiter + ${junit.version} + + + com.fasterxml.jackson.core + jackson-databind + ${jackson.version} + + + com.h2database + h2 + ${h2.version} + + + org.postgresql + postgresql + ${postgresql.version} + + + + + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.5.2 + + + + +