feat: pin running instances to immutable published definition versions

Replace register/replace with publish so new graphs can ship without rewriting old ones, and keep in-flight work on the version it started with.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-16 10:30:14 +08:00
co-authored by Cursor
parent 8bc628613d
commit 9bda7c417c
35 changed files with 688 additions and 382 deletions
+3 -3
View File
@@ -10,7 +10,7 @@
|---|---| |---|---|
| `ordo-api` | 公共模型与 `OrdoEngine` 端口 | | `ordo-api` | 公共模型与 `OrdoEngine` 端口 |
| `ordo-core` | 运行时(`DefaultOrdoEngine` / `InMemoryOrdoEngine`) | | `ordo-core` | 运行时(`DefaultOrdoEngine` / `InMemoryOrdoEngine`) |
| `ordo-storage-jdbc` | JDBC 存储 + Flyway 迁移(V1–V6) | | `ordo-storage-jdbc` | JDBC 存储 + Flyway 迁移(V1–V7) |
| `ordo-spring-boot-starter` | Spring Boot 自动装配(JDBC + Flyway) | | `ordo-spring-boot-starter` | Spring Boot 自动装配(JDBC + Flyway) |
| `ordo-example` | 内存引擎示例 | | `ordo-example` | 内存引擎示例 |
@@ -29,7 +29,7 @@
- 审计时间线:`ProcessEvent` + `queryHistory` - 审计时间线:`ProcessEvent` + `queryHistory`
- 扩展点:`AssigneeResolver`、`RoutingCondition`、`ActionHandler`、`OrdoEventListener` - 扩展点:`AssigneeResolver`、`RoutingCondition`、`ActionHandler`、`OrdoEventListener`
开发计划:定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST、目录与多版本定义之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。 开发计划:MySQL 方言、可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST 与目录之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。
详细用法(定义 JSON、扩展点、异常、查询、ACTION/审计语义)见 **[docs/usage.md](docs/usage.md)**。对外行为变更时同步更新该文档。 详细用法(定义 JSON、扩展点、异常、查询、ACTION/审计语义)见 **[docs/usage.md](docs/usage.md)**。对外行为变更时同步更新该文档。
@@ -45,7 +45,7 @@
```java ```java
OrdoEngine ordo = new InMemoryOrdoEngine(); OrdoEngine ordo = new InMemoryOrdoEngine();
ordo.register(ProcessDefinition.linear("leave-request", "Leave request", List.of( ordo.publish(ProcessDefinition.linear("leave-request", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
+3 -6
View File
@@ -7,7 +7,8 @@
## 已完成 ## 已完成
- ANY/ALL 会签/或签 - ANY/ALL 会签/或签
- JDBC 存储(PostgreSQL)+ Flyway V1–V6 - JDBC 存储(PostgreSQL)+ Flyway V1–V7
- 定义不可变多版本:`publish`;实例锁定 `definitionVersion`
- 条件路由 `StepTransition` + `RoutingCondition` - 条件路由 `StepTransition` + `RoutingCondition`
- ACTION + `ActionHandler`;执行记录持久化 - ACTION + `ActionHandler`;执行记录持久化
- 发起人撤回 `WITHDRAWN` - 发起人撤回 `WITHDRAWN`
@@ -21,10 +22,6 @@
下列能力已纳入计划,尚未实现。实现顺序可按依赖调整,但范围本身不从计划中拿掉。 下列能力已纳入计划,尚未实现。实现顺序可按依赖调整,但范围本身不从计划中拿掉。
### 定义不可变多版本
当前同 `id` 用 `replace` 整体替换;存在 `RUNNING` 实例时拒绝。计划改为定义不可变多版本:新版本不改写旧版本;运行中实例锁定发起时所用版本。
### MySQL 方言 ### MySQL 方言
当前 JDBC DDL/冲突处理面向 PostgreSQL。计划增加 MySQL 方言(及对应测试),通过 dialect 层扩展,而不是只支持一种库。 当前 JDBC DDL/冲突处理面向 PostgreSQL。计划增加 MySQL 方言(及对应测试),通过 dialect 层扩展,而不是只支持一种库。
@@ -41,7 +38,7 @@
流程设计器是**单独产品**,消费上述 REST 与目录,不做成 ordo 模块。画布对齐引擎图(审批步、ACTION 步、边上的 `when`/`priority`,结束为 `to: null`),不引入 BPMN 网关/并行等引擎没有的语义。节点坐标等 layout 由设计器自存,不进入 `ProcessDefinition`。 流程设计器是**单独产品**,消费上述 REST 与目录,不做成 ordo 模块。画布对齐引擎图(审批步、ACTION 步、边上的 `when`/`priority`,结束为 `to: null`),不引入 BPMN 网关/并行等引擎没有的语义。节点坐标等 layout 由设计器自存,不进入 `ProcessDefinition`。
启动时机:REST 契约、目录 SPI、定义不可变多版本落地之后。当前 `replace` + 无版本锁定不适合作为设计器保存/发布模型。 启动时机:REST 契约与目录 SPI 落地之后。草稿由设计器/REST 文档存储,不进入引擎图版本。
## 暂不在计划中 ## 暂不在计划中
+15 -17
View File
@@ -12,7 +12,7 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。
产品边界:不做业务表单、用户体系、多租户;不内置设计器 UI。业务字段放在 `ProcessContext`(不可变 `Map<String, Object>`)。当前也**没有** REST;HTTP 仍由宿主自建。计划中的可选 REST 与独立设计器见 [roadmap.md](roadmap.md)。 产品边界:不做业务表单、用户体系、多租户;不内置设计器 UI。业务字段放在 `ProcessContext`(不可变 `Map<String, Object>`)。当前也**没有** REST;HTTP 仍由宿主自建。计划中的可选 REST 与独立设计器见 [roadmap.md](roadmap.md)。
开发计划(尚未提供,见 [roadmap.md](roadmap.md)):定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。 开发计划(尚未提供,见 [roadmap.md](roadmap.md)):MySQL 方言、可选 REST + 目录 SPI。
## 2. 模块与接入 ## 2. 模块与接入
@@ -61,12 +61,12 @@ OrdoEngine ordo = new InMemoryOrdoEngine();
ordo: ordo:
enabled: true enabled: true
definitions: definitions:
location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 replace location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 publish
due: due:
poll-ms: 0 # >0 时轮询 processDue;默认不调度 poll-ms: 0 # >0 时轮询 processDue;默认不调度
``` ```
启动加载使用 `replace`:无 `RUNNING` 实例则整图替换;有运行中实例则保留库里的定义。 启动加载使用 `publish`:图与 latest 相同则不升版本;不同则写入新版本。运行中实例继续锁定发起时所用版本。
宿主用 `@Bean` 覆盖默认扩展点: 宿主用 `@Bean` 覆盖默认扩展点:
@@ -81,14 +81,14 @@ ordo:
## 3. 流程定义 ## 3. 流程定义
每个定义有 `id`、`name`、步骤列表、转移列表。步骤 id 在定义内唯一。每个步骤必须至少有一条出边(结束用 `to = null`)。同一 `from` 上 `priority` 不能重复。 每个定义有 `id`、引擎分配的 `version`、`name`、步骤列表、转移列表。步骤 id 在定义内唯一。每个步骤必须至少有一条出边(结束用 `to = null`)。同一 `from` 上 `priority` 不能重复。
### 3.1 代码构建 ### 3.1 代码构建
线性(每步无条件进下一步,最后一步结束): 线性(每步无条件进下一步,最后一步结束):
```java ```java
ordo.register(ProcessDefinition.linear("leave-request", "Leave request", List.of( ordo.publish(ProcessDefinition.linear("leave-request", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -113,10 +113,9 @@ new ProcessDefinition("leave-request-routed", "Leave request",
)); ));
``` ```
`register`:同 id 已存在则 `DefinitionAlreadyExistsException`。 `publish`:该 `id` 尚无版本则写入 v1;与 latest 的 id/name/steps/transitions 相同则返回 latest 不插入;否则插入 `latest + 1`。运行中实例不阻止发布。调用方构造的 `ProcessDefinition` 版本为 0;入库后由引擎分配从 1 起的单调版本。JSON 不要写 `version`,出现则忽略。
`replace`:整图覆盖;存在该定义的 `RUNNING` 实例则 `DefinitionInUseException`。
当前**没有**运行中实例锁定所用版本:`replace` 成功后新实例用新图,旧已结束实例仍按当时落库的任务理解历史。 `start(definitionId)` 使用 latest。实例带 `definitionVersion`;审批、到期、撤回、取消均按该版本取图,不跟随后续 `publish`。
### 3.2 JSON ### 3.2 JSON
@@ -303,19 +302,20 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
|---|---|---| |---|---|---|
| `queryTasks(TaskQuery, PageRequest)` | assignee、instanceId、definitionId、status、createdFrom/To | 创建时间降序 | | `queryTasks(TaskQuery, PageRequest)` | assignee、instanceId、definitionId、status、createdFrom/To | 创建时间降序 |
| `queryInstances(InstanceQuery, PageRequest)` | definitionId、status、initiator、startedFrom/To | 开始时间降序 | | `queryInstances(InstanceQuery, PageRequest)` | definitionId、status、initiator、startedFrom/To | 开始时间降序 |
| `queryDefinitions(PageRequest)` | 无过滤 | 定义 id 升序 | | `queryDefinitions(PageRequest)` | 每个 id 的 latest | 定义 id 升序 |
| `queryDefinitionVersions(id, PageRequest)` | 单 id 全部版本 | version 降序 |
| `queryHistory(instanceId, PageRequest)` | 单实例 | 发生时间升序 | | `queryHistory(instanceId, PageRequest)` | 单实例 | 发生时间升序 |
| `queryActionExecutions(instanceId, PageRequest)` | 单实例 | 开始时间升序 | | `queryActionExecutions(instanceId, PageRequest)` | 单实例 | 开始时间升序 |
`TaskQuery.any().withAssignee("maria").withStatus(TaskStatus.PENDING)` 等 with 方法返回新对象。字段 `null` 表示不按该维过滤。 `TaskQuery.any().withAssignee("maria").withStatus(TaskStatus.PENDING)` 等 with 方法返回新对象。字段 `null` 表示不按该维过滤。
便捷方法(不分页):`findInstance`、`findTask`、`findTasks`、`findPendingTasksByAssignee`、`findPendingTasksByInstanceId`。 便捷方法(不分页):`findInstance`、`findTask`、`findTasks`、`findPendingTasksByAssignee`、`findPendingTasksByInstanceId`、`findDefinition(id)` / `findDefinition(id, version)`。
## 10. `OrdoEngine` 一览 ## 10. `OrdoEngine` 一览
| 方法 | 说明 | | 方法 | 说明 |
|---|---| |---|---|
| `register` / `replace` | 登记 / 整图替换定义 | | `publish` | 发布不可变图版本(相等则 no-op) |
| `start` | 发起;可选 `ProcessContext` | | `start` | 发起;可选 `ProcessContext` |
| `approve` / `reject` | 办理当前 PENDING 任务 | | `approve` / `reject` | 办理当前 PENDING 任务 |
| `reassign` | 当前办理人把 PENDING 任务转给他人 | | `reassign` | 当前办理人把 PENDING 任务转给他人 |
@@ -323,7 +323,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| `withdraw` | 发起人撤回 | | `withdraw` | 发起人撤回 |
| `cancel` | 管理员/系统取消(引擎不鉴权角色) | | `cancel` | 管理员/系统取消(引擎不鉴权角色) |
| `find*` | 按 id / 待办索引读取 | | `find*` | 按 id / 待办索引读取 |
| `queryTasks` / `queryInstances` / `queryDefinitions` | 分页列表 | | `queryTasks` / `queryInstances` / `queryDefinitions` / `queryDefinitionVersions` | 分页列表 |
| `queryHistory` / `queryActionExecutions` | 实例审计与 ACTION 记录 | | `queryHistory` / `queryActionExecutions` | 实例审计与 ACTION 记录 |
## 11. 异常 ## 11. 异常
@@ -332,9 +332,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| 类型 | 何时 | | 类型 | 何时 |
|---|---| |---|---|
| `DefinitionAlreadyExistsException` | `register` 撞 id | | `DefinitionNotFoundException` | `start` 等找不到 latest,或实例锁定的 version 不存在 |
| `DefinitionNotFoundException` | `start` 等找不到定义 |
| `DefinitionInUseException` | `replace` 时仍有 RUNNING 实例 |
| `InstanceNotFoundException` | 撤回/取消等找不到实例 | | `InstanceNotFoundException` | 撤回/取消等找不到实例 |
| `InstanceAlreadyCompletedException` | 对非 RUNNING 实例完成/撤回/取消 | | `InstanceAlreadyCompletedException` | 对非 RUNNING 实例完成/撤回/取消 |
| `TaskNotFoundException` | 任务 id 不存在 | | `TaskNotFoundException` | 任务 id 不存在 |
@@ -347,12 +345,12 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
## 12. 存储 ## 12. 存储
Flyway 脚本在 `ordo-storage-jdbc` 的 `db/migration`(V1–V6)。表包括定义/步骤/候选人/转移、实例、任务、`ordo_process_event`、`ordo_action_execution`。 Flyway 脚本在 `ordo-storage-jdbc` 的 `db/migration`(V1–V7)。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤/候选人/转移、实例(含 `definition_version`)、任务、`ordo_process_event`、`ordo_action_execution`。
多 JVM 共享同一库时,多步写入走 `TransactionExecutor`,完成任务/实例用条件更新(仍 PENDING / 仍 RUNNING 才改),避免双花。 多 JVM 共享同一库时,多步写入走 `TransactionExecutor`,完成任务/实例用条件更新(仍 PENDING / 仍 RUNNING 才改),避免双花。
## 13. 未提供能力 ## 13. 未提供能力
开发计划中(见 [roadmap.md](roadmap.md)):定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。独立设计器不进本仓库,等 REST、目录与多版本定义之后再做。 开发计划中(见 [roadmap.md](roadmap.md)):MySQL 方言、可选 REST + 目录 SPI。独立设计器不进本仓库,等 REST 与目录之后再做。
暂不在计划中:多租户、设计器 UI。当前 REST 由宿主自建。 暂不在计划中:多租户、设计器 UI。当前 REST 由宿主自建。
@@ -8,15 +8,15 @@ import com.jetlumen.ordo.api.query.TaskQuery;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
/** Public entry point for definition registration and approval operations. */ /** Public entry point for definition publication and approval operations. */
public interface OrdoEngine { public interface OrdoEngine {
void register(ProcessDefinition definition);
/** /**
* Inserts the definition, or replaces its entire graph when the id already exists. * Publishes an immutable graph version. First publication of an id is version 1.
* Refuses replacement while any instance of this definition is {@code RUNNING}. * If the graph matches the latest version, returns that version without inserting.
* Otherwise inserts {@code latest + 1}. Running instances do not block publication.
*/ */
void replace(ProcessDefinition definition); ProcessDefinition publish(ProcessDefinition definition);
default ProcessInstance start(String definitionId, String initiator) { default ProcessInstance start(String definitionId, String initiator) {
return start(definitionId, initiator, ProcessContext.empty()); return start(definitionId, initiator, ProcessContext.empty());
} }
@@ -49,15 +49,22 @@ public interface OrdoEngine {
List<ApprovalTask> findPendingTasksByAssignee(String assignee); List<ApprovalTask> findPendingTasksByAssignee(String assignee);
List<ApprovalTask> findPendingTasksByInstanceId(String instanceId); List<ApprovalTask> findPendingTasksByInstanceId(String instanceId);
Optional<ProcessDefinition> findDefinition(String definitionId);
Optional<ProcessDefinition> findDefinition(String definitionId, int version);
/** Paginated, filterable task query; see {@link com.jetlumen.ordo.api.repository.ApprovalTaskRepository#query}. */ /** Paginated, filterable task query; see {@link com.jetlumen.ordo.api.repository.ApprovalTaskRepository#query}. */
Page<ApprovalTask> queryTasks(TaskQuery query, PageRequest pageRequest); Page<ApprovalTask> queryTasks(TaskQuery query, PageRequest pageRequest);
/** Paginated, filterable instance query; see {@link com.jetlumen.ordo.api.repository.ProcessInstanceRepository#query}. */ /** Paginated, filterable instance query; see {@link com.jetlumen.ordo.api.repository.ProcessInstanceRepository#query}. */
Page<ProcessInstance> queryInstances(InstanceQuery query, PageRequest pageRequest); Page<ProcessInstance> queryInstances(InstanceQuery query, PageRequest pageRequest);
/** Paginated listing of all registered process definitions. */ /** Paginated listing of the latest version of each process definition, ordered by id. */
Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest); Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest);
/** Paginated versions of one definition, newest version first. */
Page<ProcessDefinition> queryDefinitionVersions(String definitionId, PageRequest pageRequest);
/** Instance timeline, oldest-first; see {@link com.jetlumen.ordo.api.repository.ProcessHistoryRepository#query}. */ /** Instance timeline, oldest-first; see {@link com.jetlumen.ordo.api.repository.ProcessHistoryRepository#query}. */
Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest); Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest);
@@ -8,9 +8,13 @@ import java.util.Objects;
import java.util.Set; import java.util.Set;
/** Immutable blueprint for an approval process with explicit step transitions. */ /** Immutable blueprint for an approval process with explicit step transitions. */
public record ProcessDefinition(String id, String name, List<ApprovalStep> steps, List<StepTransition> transitions) { public record ProcessDefinition(String id, int version, String name, List<ApprovalStep> steps,
List<StepTransition> transitions) {
public ProcessDefinition { public ProcessDefinition {
ApprovalStep.requireText(id, "definition id"); ApprovalStep.requireText(id, "definition id");
if (version < 0) {
throw new IllegalArgumentException("definition version must not be negative");
}
ApprovalStep.requireText(name, "definition name"); ApprovalStep.requireText(name, "definition name");
steps = List.copyOf(steps); steps = List.copyOf(steps);
if (steps.isEmpty()) { if (steps.isEmpty()) {
@@ -52,6 +56,24 @@ public record ProcessDefinition(String id, String name, List<ApprovalStep> steps
.toList(); .toList();
} }
/** Unpublished graph; the engine assigns a version on {@code publish}. */
public ProcessDefinition(String id, String name, List<ApprovalStep> steps, List<StepTransition> transitions) {
this(id, 0, name, steps, transitions);
}
public ProcessDefinition withVersion(int version) {
return new ProcessDefinition(id, version, name, steps, transitions);
}
/** Equality of the executable graph, ignoring assigned version. */
public boolean sameGraph(ProcessDefinition other) {
Objects.requireNonNull(other, "other must not be null");
return id.equals(other.id)
&& name.equals(other.name)
&& steps.equals(other.steps)
&& transitions.equals(other.transitions);
}
/** /**
* Builds a definition whose transitions mirror the former linear steps order: each step * Builds a definition whose transitions mirror the former linear steps order: each step
* unconditionally advances to the next, and the last step unconditionally ends. * unconditionally advances to the next, and the last step unconditionally ends.
@@ -112,6 +112,7 @@ public final class ProcessDefinitionParser {
private record DefinitionDocument( private record DefinitionDocument(
String id, String id,
String name, String name,
Integer version,
String startStep, String startStep,
List<StepDocument> steps, List<StepDocument> steps,
List<TransitionDocument> transitions) { List<TransitionDocument> transitions) {
@@ -2,6 +2,6 @@ package com.jetlumen.ordo.api;
import java.time.Instant; import java.time.Instant;
public record ProcessInstance(String id, String definitionId, String initiator, ProcessStatus status, public record ProcessInstance(String id, String definitionId, int definitionVersion, String initiator,
Instant startedAt, Instant finishedAt, ProcessContext context) { ProcessStatus status, Instant startedAt, Instant finishedAt, ProcessContext context) {
} }
@@ -1,7 +0,0 @@
package com.jetlumen.ordo.api.exception;
public final class DefinitionAlreadyExistsException extends OrdoException {
public DefinitionAlreadyExistsException(String definitionId) {
super("definition already exists: " + definitionId);
}
}
@@ -1,7 +0,0 @@
package com.jetlumen.ordo.api.exception;
public final class DefinitionInUseException extends OrdoException {
public DefinitionInUseException(String definitionId) {
super("definition in use: " + definitionId);
}
}
@@ -4,4 +4,8 @@ public final class DefinitionNotFoundException extends OrdoException {
public DefinitionNotFoundException(String definitionId) { public DefinitionNotFoundException(String definitionId) {
super("definition not found: " + definitionId); super("definition not found: " + definitionId);
} }
public DefinitionNotFoundException(String definitionId, int version) {
super("definition not found: " + definitionId + " version " + version);
}
} }
@@ -6,23 +6,21 @@ import com.jetlumen.ordo.api.query.PageRequest;
import java.util.Optional; import java.util.Optional;
/** Storage port for process definitions. */ /** Storage port for immutable process definition versions. */
public interface ProcessDefinitionRepository { public interface ProcessDefinitionRepository {
/** /**
* Inserts the definition if no definition with the same id exists. * Inserts version 1 when the id is new, returns the latest version when the graph is unchanged,
* * otherwise inserts {@code latest + 1}.
* @return true if the definition was inserted, false if a definition with the same id already exists
*/ */
boolean insertIfAbsent(ProcessDefinition definition); ProcessDefinition publish(ProcessDefinition definition);
/** Optional<ProcessDefinition> findLatest(String definitionId);
* Inserts the definition, or replaces its name, steps, candidates and transitions
* when the id already exists.
*/
void upsert(ProcessDefinition definition);
Optional<ProcessDefinition> findById(String definitionId); Optional<ProcessDefinition> find(String definitionId, int version);
/** Paginated listing of all registered definitions, ordered by id ascending. */ /** Paginated listing of the latest version of each definition, ordered by id ascending. */
Page<ProcessDefinition> findAll(PageRequest pageRequest); Page<ProcessDefinition> findAll(PageRequest pageRequest);
/** Paginated versions of one definition, ordered by version descending. */
Page<ProcessDefinition> findVersions(String definitionId, PageRequest pageRequest);
} }
@@ -38,6 +38,25 @@ class ProcessDefinitionParserTest {
)); ));
assertEquals(expected, ProcessDefinitionParser.fromJson(LEAVE_REQUEST_JSON)); assertEquals(expected, ProcessDefinitionParser.fromJson(LEAVE_REQUEST_JSON));
assertEquals(0, ProcessDefinitionParser.fromJson(LEAVE_REQUEST_JSON).version());
}
@Test
void ignoresVersionInJson() {
String json = """
{
"id": "leave",
"name": "Leave request",
"version": 9,
"steps": [
{ "id": "manager", "name": "Manager approval", "candidates": ["maria"] }
],
"transitions": [
{ "from": "manager", "to": null }
]
}
""";
assertEquals(0, ProcessDefinitionParser.fromJson(json).version());
} }
@Test @Test
@@ -22,8 +22,6 @@ import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.TaskAction; import com.jetlumen.ordo.api.TaskAction;
import com.jetlumen.ordo.api.TaskStatus; import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.api.TransactionExecutor; import com.jetlumen.ordo.api.TransactionExecutor;
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException; import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.InstanceNotFoundException; import com.jetlumen.ordo.api.exception.InstanceNotFoundException;
@@ -102,26 +100,9 @@ public final class DefaultOrdoEngine implements OrdoEngine {
} }
@Override @Override
public synchronized void register(ProcessDefinition definition) { public synchronized ProcessDefinition publish(ProcessDefinition definition) {
Objects.requireNonNull(definition, "definition must not be null"); Objects.requireNonNull(definition, "definition must not be null");
transactionExecutor.execute(() -> { return transactionExecutor.execute(() -> definitionRepository.publish(definition));
if (!definitionRepository.insertIfAbsent(definition)) {
throw new DefinitionAlreadyExistsException(definition.id());
}
return null;
});
}
@Override
public synchronized void replace(ProcessDefinition definition) {
Objects.requireNonNull(definition, "definition must not be null");
transactionExecutor.execute(() -> {
if (instanceRepository.existsRunning(definition.id())) {
throw new DefinitionInUseException(definition.id());
}
definitionRepository.upsert(definition);
return null;
});
} }
@Override @Override
@@ -131,9 +112,9 @@ public final class DefaultOrdoEngine implements OrdoEngine {
List<PendingAction> queued = new ArrayList<>(); List<PendingAction> queued = new ArrayList<>();
List<ProcessEvent> events = new ArrayList<>(); List<ProcessEvent> events = new ArrayList<>();
ProcessInstance instance = transactionExecutor.execute(() -> { ProcessInstance instance = transactionExecutor.execute(() -> {
ProcessDefinition definition = requireDefinition(definitionId); ProcessDefinition definition = requireLatestDefinition(definitionId);
Instant now = clock.instant(); Instant now = clock.instant();
ProcessInstance started = new ProcessInstance(nextId(), definition.id(), initiator, ProcessInstance started = new ProcessInstance(nextId(), definition.id(), definition.version(), initiator,
ProcessStatus.RUNNING, now, null, context); ProcessStatus.RUNNING, now, null, context);
instanceRepository.insert(started); instanceRepository.insert(started);
record(events, started.id(), null, null, ProcessEventType.INSTANCE_STARTED, initiator, null, now); record(events, started.id(), null, null, ProcessEventType.INSTANCE_STARTED, initiator, null, now);
@@ -232,8 +213,9 @@ public final class DefaultOrdoEngine implements OrdoEngine {
throw new InstanceAlreadyCompletedException(instanceId); throw new InstanceAlreadyCompletedException(instanceId);
} }
Instant now = clock.instant(); Instant now = clock.instant();
ProcessInstance finished = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(), ProcessInstance finished = new ProcessInstance(instance.id(), instance.definitionId(),
status, instance.startedAt(), now, instance.context()); instance.definitionVersion(), instance.initiator(), status, instance.startedAt(), now,
instance.context());
if (!instanceRepository.completeIfRunning(finished)) { if (!instanceRepository.completeIfRunning(finished)) {
throw new InstanceAlreadyCompletedException(instanceId); throw new InstanceAlreadyCompletedException(instanceId);
} }
@@ -295,12 +277,34 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return instanceRepository.query(query, pageRequest); return instanceRepository.query(query, pageRequest);
} }
@Override
public synchronized Optional<ProcessDefinition> findDefinition(String definitionId) {
requireText(definitionId, "definition id");
return definitionRepository.findLatest(definitionId);
}
@Override
public synchronized Optional<ProcessDefinition> findDefinition(String definitionId, int version) {
requireText(definitionId, "definition id");
if (version < 1) {
throw new IllegalArgumentException("definition version must be positive");
}
return definitionRepository.find(definitionId, version);
}
@Override @Override
public synchronized Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest) { public synchronized Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest) {
Objects.requireNonNull(pageRequest, "pageRequest must not be null"); Objects.requireNonNull(pageRequest, "pageRequest must not be null");
return definitionRepository.findAll(pageRequest); return definitionRepository.findAll(pageRequest);
} }
@Override
public synchronized Page<ProcessDefinition> queryDefinitionVersions(String definitionId, PageRequest pageRequest) {
requireText(definitionId, "definition id");
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
return definitionRepository.findVersions(definitionId, pageRequest);
}
@Override @Override
public synchronized Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) { public synchronized Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) {
requireText(instanceId, "instance id"); requireText(instanceId, "instance id");
@@ -345,7 +349,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
if (instance.status() != ProcessStatus.RUNNING) { if (instance.status() != ProcessStatus.RUNNING) {
return false; return false;
} }
ProcessDefinition definition = requireDefinition(instance.definitionId()); ProcessDefinition definition = requireDefinition(instance);
ApprovalStep step = requireStep(definition, overdue.stepId()); ApprovalStep step = requireStep(definition, overdue.stepId());
StepDue due = step.due(); StepDue due = step.due();
if (due == null) { if (due == null) {
@@ -454,7 +458,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
private void advanceAfterDecision(ApprovalTask completedTask, Instant now, List<PendingAction> queued, private void advanceAfterDecision(ApprovalTask completedTask, Instant now, List<PendingAction> queued,
List<ProcessEvent> events) { List<ProcessEvent> events) {
ProcessInstance instance = requireInstance(completedTask.instanceId()); ProcessInstance instance = requireInstance(completedTask.instanceId());
ProcessDefinition definition = requireDefinition(instance.definitionId()); ProcessDefinition definition = requireDefinition(instance);
ApprovalStep step = requireStep(definition, completedTask.stepId()); ApprovalStep step = requireStep(definition, completedTask.stepId());
List<ApprovalTask> siblings = taskRepository.findByInstanceIdAndStepId(instance.id(), step.id()); List<ApprovalTask> siblings = taskRepository.findByInstanceIdAndStepId(instance.id(), step.id());
@@ -600,8 +604,9 @@ public final class DefaultOrdoEngine implements OrdoEngine {
private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now, private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now,
List<ProcessEvent> events) { List<ProcessEvent> events) {
ProcessInstance completed = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(), ProcessInstance completed = new ProcessInstance(instance.id(), instance.definitionId(),
status, instance.startedAt(), now, instance.context()); instance.definitionVersion(), instance.initiator(), status, instance.startedAt(), now,
instance.context());
if (!instanceRepository.completeIfRunning(completed)) { if (!instanceRepository.completeIfRunning(completed)) {
throw new InstanceAlreadyCompletedException(instance.id()); throw new InstanceAlreadyCompletedException(instance.id());
} }
@@ -610,11 +615,17 @@ public final class DefaultOrdoEngine implements OrdoEngine {
record(events, instance.id(), null, null, type, null, null, now); record(events, instance.id(), null, null, type, null, null, now);
} }
private ProcessDefinition requireDefinition(String definitionId) { private ProcessDefinition requireLatestDefinition(String definitionId) {
return definitionRepository.findById(definitionId) return definitionRepository.findLatest(definitionId)
.orElseThrow(() -> new DefinitionNotFoundException(definitionId)); .orElseThrow(() -> new DefinitionNotFoundException(definitionId));
} }
private ProcessDefinition requireDefinition(ProcessInstance instance) {
return definitionRepository.find(instance.definitionId(), instance.definitionVersion())
.orElseThrow(() -> new DefinitionNotFoundException(instance.definitionId(),
instance.definitionVersion()));
}
private ProcessInstance requireInstance(String instanceId) { private ProcessInstance requireInstance(String instanceId) {
return instanceRepository.findById(instanceId) return instanceRepository.findById(instanceId)
.orElseThrow(() -> new IllegalStateException("instance not found: " + instanceId)); .orElseThrow(() -> new IllegalStateException("instance not found: " + instanceId));
@@ -76,13 +76,8 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
} }
@Override @Override
public void register(ProcessDefinition definition) { public ProcessDefinition publish(ProcessDefinition definition) {
delegate.register(definition); return delegate.publish(definition);
}
@Override
public void replace(ProcessDefinition definition) {
delegate.replace(definition);
} }
@Override @Override
@@ -150,11 +145,26 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
return delegate.queryInstances(query, pageRequest); return delegate.queryInstances(query, pageRequest);
} }
@Override
public Optional<ProcessDefinition> findDefinition(String definitionId) {
return delegate.findDefinition(definitionId);
}
@Override
public Optional<ProcessDefinition> findDefinition(String definitionId, int version) {
return delegate.findDefinition(definitionId, version);
}
@Override @Override
public Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest) { public Page<ProcessDefinition> queryDefinitions(PageRequest pageRequest) {
return delegate.queryDefinitions(pageRequest); return delegate.queryDefinitions(pageRequest);
} }
@Override
public Page<ProcessDefinition> queryDefinitionVersions(String definitionId, PageRequest pageRequest) {
return delegate.queryDefinitionVersions(definitionId, pageRequest);
}
@Override @Override
public Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) { public Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) {
return delegate.queryHistory(instanceId, pageRequest); return delegate.queryHistory(instanceId, pageRequest);
@@ -5,35 +5,63 @@ import com.jetlumen.ordo.api.query.Page;
import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.query.PageRequest;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import java.util.ArrayList;
import java.util.Comparator; import java.util.Comparator;
import java.util.HashMap; import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Objects;
import java.util.Optional; import java.util.Optional;
/** Development-only in-memory implementation of the definition storage port. */ /** Development-only in-memory implementation of the definition storage port. */
public final class InMemoryProcessDefinitionRepository implements ProcessDefinitionRepository { public final class InMemoryProcessDefinitionRepository implements ProcessDefinitionRepository {
private final Map<String, ProcessDefinition> definitions = new HashMap<>(); private final Map<String, Integer> latestVersions = new HashMap<>();
private final Map<String, Map<Integer, ProcessDefinition>> versions = new HashMap<>();
@Override @Override
public synchronized boolean insertIfAbsent(ProcessDefinition definition) { public synchronized ProcessDefinition publish(ProcessDefinition definition) {
return definitions.putIfAbsent(definition.id(), definition) == null; Objects.requireNonNull(definition, "definition must not be null");
Integer current = latestVersions.get(definition.id());
if (current == null) {
ProcessDefinition first = definition.withVersion(1);
versions.computeIfAbsent(definition.id(), id -> new HashMap<>()).put(1, first);
latestVersions.put(definition.id(), 1);
return first;
}
ProcessDefinition latest = versions.get(definition.id()).get(current);
if (latest.sameGraph(definition)) {
return latest;
}
int next = current + 1;
ProcessDefinition published = definition.withVersion(next);
versions.get(definition.id()).put(next, published);
latestVersions.put(definition.id(), next);
return published;
} }
@Override @Override
public synchronized void upsert(ProcessDefinition definition) { public synchronized Optional<ProcessDefinition> findLatest(String definitionId) {
definitions.put(definition.id(), definition); Integer version = latestVersions.get(definitionId);
if (version == null) {
return Optional.empty();
}
return Optional.of(versions.get(definitionId).get(version));
} }
@Override @Override
public synchronized Optional<ProcessDefinition> findById(String definitionId) { public synchronized Optional<ProcessDefinition> find(String definitionId, int version) {
return Optional.ofNullable(definitions.get(definitionId)); Map<Integer, ProcessDefinition> byVersion = versions.get(definitionId);
if (byVersion == null) {
return Optional.empty();
}
return Optional.ofNullable(byVersion.get(version));
} }
@Override @Override
public synchronized Page<ProcessDefinition> findAll(PageRequest pageRequest) { public synchronized Page<ProcessDefinition> findAll(PageRequest pageRequest) {
List<ProcessDefinition> sorted = definitions.values().stream() List<ProcessDefinition> sorted = latestVersions.keySet().stream()
.sorted(Comparator.comparing(ProcessDefinition::id)) .sorted(Comparator.naturalOrder())
.map(id -> versions.get(id).get(latestVersions.get(id)))
.toList(); .toList();
List<ProcessDefinition> page = sorted.stream() List<ProcessDefinition> page = sorted.stream()
.skip((long) pageRequest.offset()) .skip((long) pageRequest.offset())
@@ -41,4 +69,16 @@ public final class InMemoryProcessDefinitionRepository implements ProcessDefinit
.toList(); .toList();
return new Page<>(page, sorted.size(), pageRequest.page(), pageRequest.size()); return new Page<>(page, sorted.size(), pageRequest.page(), pageRequest.size());
} }
@Override
public synchronized Page<ProcessDefinition> findVersions(String definitionId, PageRequest pageRequest) {
Map<Integer, ProcessDefinition> byVersion = versions.getOrDefault(definitionId, Map.of());
List<ProcessDefinition> sorted = new ArrayList<>(byVersion.values());
sorted.sort(Comparator.comparingInt(ProcessDefinition::version).reversed());
List<ProcessDefinition> page = sorted.stream()
.skip((long) pageRequest.offset())
.limit(pageRequest.size())
.toList();
return new Page<>(page, sorted.size(), pageRequest.page(), pageRequest.size());
}
} }
@@ -17,8 +17,6 @@ import com.jetlumen.ordo.api.StepDue;
import com.jetlumen.ordo.api.StepKind; import com.jetlumen.ordo.api.StepKind;
import com.jetlumen.ordo.api.StepTransition; import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.TaskStatus; import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException; import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.InstanceNotFoundException; import com.jetlumen.ordo.api.exception.InstanceNotFoundException;
@@ -53,7 +51,7 @@ class InMemoryOrdoEngineTest {
@BeforeEach @BeforeEach
void setUp() { void setUp() {
engine = new InMemoryOrdoEngine(); engine = new InMemoryOrdoEngine();
engine.register(ProcessDefinition.linear("leave", "Leave request", List.of( engine.publish(ProcessDefinition.linear("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -136,7 +134,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void reassignDoesNotAdvanceAnAllPolicyStepAndRejectsDuplicatePendingAssignee() { void reassignDoesNotAdvanceAnAllPolicyStepAndRejectsDuplicatePendingAssignee() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(ProcessDefinition.linear("leave-all-reassign", "Leave request", List.of( allEngine.publish(ProcessDefinition.linear("leave-all-reassign", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL) new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL)
))); )));
var instance = allEngine.start("leave-all-reassign", "alice"); var instance = allEngine.start("leave-all-reassign", "alice");
@@ -157,51 +155,56 @@ class InMemoryOrdoEngineTest {
} }
@Test @Test
void exposesSpecificExceptionsForMissingAndDuplicateResources() { void exposesSpecificExceptionsForMissingResources() {
assertThrows(DefinitionNotFoundException.class, () -> engine.start("missing", "alice")); assertThrows(DefinitionNotFoundException.class, () -> engine.start("missing", "alice"));
assertThrows(TaskNotFoundException.class, () -> engine.approve("missing", "maria")); assertThrows(TaskNotFoundException.class, () -> engine.approve("missing", "maria"));
assertThrows(DefinitionAlreadyExistsException.class, () -> engine.register(ProcessDefinition.linear(
"leave", "Another leave request", List.of(ApprovalStep.single("lead", "Lead approval", "lee")))));
} }
@Test @Test
void replaceInsertsWhenTheDefinitionIsMissing() { void publishInsertsVersionOneWhenTheDefinitionIsMissing() {
InMemoryOrdoEngine empty = new InMemoryOrdoEngine(); InMemoryOrdoEngine empty = new InMemoryOrdoEngine();
empty.replace(ProcessDefinition.linear("expense", "Expense request", List.of( ProcessDefinition published = empty.publish(ProcessDefinition.linear("expense", "Expense request", List.of(
ApprovalStep.single("director", "Director approval", "diana")))); ApprovalStep.single("director", "Director approval", "diana"))));
assertEquals(1, published.version());
var instance = empty.start("expense", "alice"); var instance = empty.start("expense", "alice");
assertEquals(1, instance.definitionVersion());
assertEquals("diana", empty.findTasks(instance.id()).get(0).assignee()); assertEquals("diana", empty.findTasks(instance.id()).get(0).assignee());
} }
@Test @Test
void replaceSwapsTheGraphWhenNoInstanceIsRunning() { void publishIsIdempotentWhenTheGraphIsUnchanged() {
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of( ProcessDefinition first = engine.findDefinition("leave").orElseThrow();
ApprovalStep.single("director", "Director approval", "diana")))); ProcessDefinition second = engine.publish(ProcessDefinition.linear("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"),
var instance = engine.start("leave", "alice"); ApprovalStep.single("hr", "HR approval", "henry")
ApprovalTask task = engine.findTasks(instance.id()).get(0); )));
assertEquals("director", task.stepId()); assertEquals(first, second);
assertEquals("diana", task.assignee()); assertEquals(1, engine.queryDefinitionVersions("leave", new PageRequest(0, 10)).totalElements());
} }
@Test @Test
void replaceIsRejectedWhileAnInstanceIsRunning() { void publishCreatesANewVersionAndLocksRunningInstancesToTheOldGraph() {
engine.start("leave", "alice"); var running = engine.start("leave", "alice");
assertThrows(DefinitionInUseException.class, () -> engine.replace(ProcessDefinition.linear( assertEquals(1, running.definitionVersion());
"leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana")))));
}
@Test engine.publish(ProcessDefinition.linear("leave", "Leave request v2", List.of(
void replaceSucceedsAfterInstancesReachATerminalStatus() {
var instance = engine.start("leave", "alice");
engine.reject(engine.findTasks(instance.id()).get(0).id(), "maria");
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of(
ApprovalStep.single("director", "Director approval", "diana")))); ApprovalStep.single("director", "Director approval", "diana"))));
assertEquals(2, engine.findDefinition("leave").orElseThrow().version());
assertEquals("maria", engine.findTasks(running.id()).get(0).assignee());
engine.approve(engine.findTasks(running.id()).get(0).id(), "maria");
assertEquals("hr", engine.findPendingTasksByInstanceId(running.id()).get(0).stepId());
var next = engine.start("leave", "bob"); var next = engine.start("leave", "bob");
assertEquals(2, next.definitionVersion());
assertEquals("diana", engine.findTasks(next.id()).get(0).assignee()); assertEquals("diana", engine.findTasks(next.id()).get(0).assignee());
Page<ProcessDefinition> latest = engine.queryDefinitions(new PageRequest(0, 10));
assertEquals(List.of("leave"), latest.content().stream().map(ProcessDefinition::id).toList());
assertEquals(List.of(2), latest.content().stream().map(ProcessDefinition::version).toList());
Page<ProcessDefinition> versions = engine.queryDefinitionVersions("leave", new PageRequest(0, 10));
assertEquals(List.of(2, 1), versions.content().stream().map(ProcessDefinition::version).toList());
} }
@Test @Test
@@ -352,7 +355,7 @@ class InMemoryOrdoEngineTest {
.filter(String.class::isInstance) .filter(String.class::isInstance)
.map(String.class::cast) .map(String.class::cast)
.orElse(candidate)); .orElse(candidate));
contextAwareEngine.register(ProcessDefinition.linear("leave", "Leave request", List.of( contextAwareEngine.publish(ProcessDefinition.linear("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -372,7 +375,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void anyPolicyAdvancesOnFirstApprovalAndSkipsTheOtherCandidates() { void anyPolicyAdvancesOnFirstApprovalAndSkipsTheOtherCandidates() {
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
anyEngine.register(ProcessDefinition.linear("leave-any", "Leave request", List.of( anyEngine.publish(ProcessDefinition.linear("leave-any", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY), new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -401,7 +404,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void anyPolicyOnlyRejectsTheStepOnceEveryCandidateHasRejected() { void anyPolicyOnlyRejectsTheStepOnceEveryCandidateHasRejected() {
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
anyEngine.register(ProcessDefinition.linear("leave-any-reject", "Leave request", List.of( anyEngine.publish(ProcessDefinition.linear("leave-any-reject", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY) new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY)
))); )));
@@ -421,7 +424,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void allPolicyOnlyAdvancesOnceEveryCandidateHasApproved() { void allPolicyOnlyAdvancesOnceEveryCandidateHasApproved() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(ProcessDefinition.linear("leave-all", "Leave request", List.of( allEngine.publish(ProcessDefinition.linear("leave-all", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL), new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -446,7 +449,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void allPolicyFailsFastAndSkipsRemainingCandidatesOnASingleRejection() { void allPolicyFailsFastAndSkipsRemainingCandidatesOnASingleRejection() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(ProcessDefinition.linear("leave-all-reject", "Leave request", List.of( allEngine.publish(ProcessDefinition.linear("leave-all-reject", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL) new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL)
))); )));
@@ -478,7 +481,7 @@ class InMemoryOrdoEngineTest {
List<ApprovalStep> steps = List.of( List<ApprovalStep> steps = List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("director", "Director approval", "diana")); ApprovalStep.single("director", "Director approval", "diana"));
routingEngine.register(new ProcessDefinition("expense", "Expense request", steps, List.of( routingEngine.publish(new ProcessDefinition("expense", "Expense request", steps, List.of(
StepTransition.when("manager", "director", "amount-gt-1000", 0), StepTransition.when("manager", "director", "amount-gt-1000", 0),
new StepTransition("manager", null, null, 1), new StepTransition("manager", null, null, 1),
StepTransition.end("director")))); StepTransition.end("director"))));
@@ -500,7 +503,7 @@ class InMemoryOrdoEngineTest {
List<String> executed = new java.util.ArrayList<>(); List<String> executed = new java.util.ArrayList<>();
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(), InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
RoutingCondition.always(), (key, context) -> executed.add(key)); RoutingCondition.always(), (key, context) -> executed.add(key));
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of( actionEngine.publish(new ProcessDefinition("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.action("notify", "Notify HR", "leave-approved-mail"), ApprovalStep.action("notify", "Notify HR", "leave-approved-mail"),
ApprovalStep.single("hr", "HR approval", "henry")), ApprovalStep.single("hr", "HR approval", "henry")),
@@ -524,7 +527,7 @@ class InMemoryOrdoEngineTest {
RoutingCondition.always(), (key, context) -> { RoutingCondition.always(), (key, context) -> {
throw new IllegalStateException("mail failed"); throw new IllegalStateException("mail failed");
}); });
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of( actionEngine.publish(new ProcessDefinition("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.action("notify", "Notify HR", "leave-approved-mail")), ApprovalStep.action("notify", "Notify HR", "leave-approved-mail")),
List.of( List.of(
@@ -541,7 +544,7 @@ class InMemoryOrdoEngineTest {
List<String> executed = new java.util.ArrayList<>(); List<String> executed = new java.util.ArrayList<>();
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(), InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
RoutingCondition.always(), (key, context) -> executed.add(key)); RoutingCondition.always(), (key, context) -> executed.add(key));
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of( actionEngine.publish(new ProcessDefinition("leave", "Leave request", List.of(
ApprovalStep.action("notify", "Notify manager", "leave-submitted-mail"), ApprovalStep.action("notify", "Notify manager", "leave-submitted-mail"),
ApprovalStep.single("manager", "Manager approval", "maria")), ApprovalStep.single("manager", "Manager approval", "maria")),
List.of( List.of(
@@ -557,7 +560,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void throwsWhenNoTransitionMatches() { void throwsWhenNoTransitionMatches() {
InMemoryOrdoEngine routingEngine = new InMemoryOrdoEngine((key, context) -> false); InMemoryOrdoEngine routingEngine = new InMemoryOrdoEngine((key, context) -> false);
routingEngine.register(new ProcessDefinition("expense", "Expense request", routingEngine.publish(new ProcessDefinition("expense", "Expense request",
List.of(ApprovalStep.single("manager", "Manager approval", "maria")), List.of(ApprovalStep.single("manager", "Manager approval", "maria")),
List.of(StepTransition.endWhen("manager", "never", 0)))); List.of(StepTransition.endWhen("manager", "never", 0))));
@@ -568,7 +571,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void queryTasksFiltersByAssigneeDefinitionAndPaginates() { void queryTasksFiltersByAssigneeDefinitionAndPaginates() {
engine.register(ProcessDefinition.linear("expense", "Expense request", List.of( engine.publish(ProcessDefinition.linear("expense", "Expense request", List.of(
ApprovalStep.single("finance", "Finance approval", "frank")))); ApprovalStep.single("finance", "Finance approval", "frank"))));
var leaveInstance = engine.start("leave", "alice"); var leaveInstance = engine.start("leave", "alice");
engine.start("expense", "alice"); engine.start("expense", "alice");
@@ -614,7 +617,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void queryDefinitionsPaginatesRegisteredDefinitions() { void queryDefinitionsPaginatesRegisteredDefinitions() {
engine.register(ProcessDefinition.linear("expense", "Expense request", List.of( engine.publish(ProcessDefinition.linear("expense", "Expense request", List.of(
ApprovalStep.single("finance", "Finance approval", "frank")))); ApprovalStep.single("finance", "Finance approval", "frank"))));
Page<ProcessDefinition> all = engine.queryDefinitions(new PageRequest(0, 10)); Page<ProcessDefinition> all = engine.queryDefinitions(new PageRequest(0, 10));
@@ -625,7 +628,7 @@ class InMemoryOrdoEngineTest {
@Test @Test
void recordsHistoryForStartApproveSkipAndComplete() { void recordsHistoryForStartApproveSkipAndComplete() {
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine(); InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
anyEngine.register(ProcessDefinition.linear("leave-any", "Leave request", List.of( anyEngine.publish(ProcessDefinition.linear("leave-any", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY), new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -690,7 +693,7 @@ class InMemoryOrdoEngineTest {
} }
received.add(event.type()); received.add(event.type());
})); }));
listening.register(ProcessDefinition.linear("leave", "Leave request", List.of( listening.publish(ProcessDefinition.linear("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria")))); ApprovalStep.single("manager", "Manager approval", "maria"))));
var instance = listening.start("leave", "alice"); var instance = listening.start("leave", "alice");
@@ -711,7 +714,7 @@ class InMemoryOrdoEngineTest {
throw new IllegalStateException("mail failed"); throw new IllegalStateException("mail failed");
} }
}); });
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of( actionEngine.publish(new ProcessDefinition("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.action("notify", "Notify HR", "ok-mail"), ApprovalStep.action("notify", "Notify HR", "ok-mail"),
ApprovalStep.action("fail", "Fail mail", "fail-mail")), ApprovalStep.action("fail", "Fail mail", "fail-mail")),
@@ -743,7 +746,7 @@ class InMemoryOrdoEngineTest {
void processDueReassignsAfterTheStepDueElapses() { void processDueReassignsAfterTheStepDueElapses() {
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z")); MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock); InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock);
dueEngine.register(ProcessDefinition.linear("leave-due", "Leave request", List.of( dueEngine.publish(ProcessDefinition.linear("leave-due", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL, new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
null, StepDue.reassign(java.time.Duration.ofHours(1), "diana"))))); null, StepDue.reassign(java.time.Duration.ofHours(1), "diana")))));
@@ -770,7 +773,7 @@ class InMemoryOrdoEngineTest {
List<String> actions = new java.util.ArrayList<>(); List<String> actions = new java.util.ArrayList<>();
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock, AssigneeResolver.direct(), InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock, AssigneeResolver.direct(),
RoutingCondition.always(), (key, context) -> actions.add(key)); RoutingCondition.always(), (key, context) -> actions.add(key));
dueEngine.register(ProcessDefinition.linear("leave-notify", "Leave request", List.of( dueEngine.publish(ProcessDefinition.linear("leave-notify", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL, new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
null, StepDue.notify(java.time.Duration.ofMinutes(30), "overdue-mail"))))); null, StepDue.notify(java.time.Duration.ofMinutes(30), "overdue-mail")))));
@@ -786,7 +789,7 @@ class InMemoryOrdoEngineTest {
void processDueGotoSkipsTheStepAndEntersTheTarget() { void processDueGotoSkipsTheStepAndEntersTheTarget() {
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z")); MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock); InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock);
dueEngine.register(new ProcessDefinition("leave-goto", "Leave request", List.of( dueEngine.publish(new ProcessDefinition("leave-goto", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL, new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
null, StepDue.gotoStep(java.time.Duration.ofHours(1), "hr")), null, StepDue.gotoStep(java.time.Duration.ofHours(1), "hr")),
ApprovalStep.single("hr", "HR approval", "henry")), ApprovalStep.single("hr", "HR approval", "henry")),
@@ -61,9 +61,9 @@ class InMemoryApprovalTaskRepositoryTest {
@Test @Test
void definitionIdFilterResolvesThroughTheInstanceRepository() { void definitionIdFilterResolvesThroughTheInstanceRepository() {
InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository(); InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository();
instanceRepository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, instanceRepository.insert(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty())); CREATED_AT, null, ProcessContext.empty()));
instanceRepository.insert(new ProcessInstance("inst-2", "expense", "alice", ProcessStatus.RUNNING, instanceRepository.insert(new ProcessInstance("inst-2", "expense", 1, "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty())); CREATED_AT, null, ProcessContext.empty()));
InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository(instanceRepository); InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository(instanceRepository);
ApprovalTask leaveTask = task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT); ApprovalTask leaveTask = task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT);
@@ -14,11 +14,11 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
class InMemoryProcessDefinitionRepositoryTest { class InMemoryProcessDefinitionRepositoryTest {
@Test @Test
void findAllPaginatesDefinitionsOrderedById() { void findAllPaginatesLatestDefinitionsOrderedById() {
InMemoryProcessDefinitionRepository repository = new InMemoryProcessDefinitionRepository(); InMemoryProcessDefinitionRepository repository = new InMemoryProcessDefinitionRepository();
repository.insertIfAbsent(definition("c-def")); repository.publish(definition("c-def"));
repository.insertIfAbsent(definition("a-def")); repository.publish(definition("a-def"));
repository.insertIfAbsent(definition("b-def")); repository.publish(definition("b-def"));
Page<ProcessDefinition> pageOne = repository.findAll(new PageRequest(0, 2)); Page<ProcessDefinition> pageOne = repository.findAll(new PageRequest(0, 2));
assertEquals(3, pageOne.totalElements()); assertEquals(3, pageOne.totalElements());
@@ -31,6 +31,20 @@ class InMemoryProcessDefinitionRepositoryTest {
assertFalse(pageTwo.hasNext()); assertFalse(pageTwo.hasNext());
} }
@Test
void publishKeepsPreviousVersions() {
InMemoryProcessDefinitionRepository repository = new InMemoryProcessDefinitionRepository();
repository.publish(ProcessDefinition.linear("leave", "v1", List.of(
ApprovalStep.single("lead", "Lead approval", "lee"))));
repository.publish(ProcessDefinition.linear("leave", "v2", List.of(
ApprovalStep.single("director", "Director approval", "diana"))));
assertEquals(2, repository.findLatest("leave").orElseThrow().version());
assertEquals("v1", repository.find("leave", 1).orElseThrow().name());
Page<ProcessDefinition> versions = repository.findVersions("leave", new PageRequest(0, 10));
assertEquals(List.of(2, 1), versions.content().stream().map(ProcessDefinition::version).toList());
}
private static ProcessDefinition definition(String id) { private static ProcessDefinition definition(String id) {
return ProcessDefinition.linear(id, id, List.of(ApprovalStep.single("lead", "Lead approval", "lee"))); return ProcessDefinition.linear(id, id, List.of(ApprovalStep.single("lead", "Lead approval", "lee")));
} }
@@ -56,6 +56,6 @@ class InMemoryProcessInstanceRepositoryTest {
private static ProcessInstance instance(String id, String definitionId, String initiator, ProcessStatus status, private static ProcessInstance instance(String id, String definitionId, String initiator, ProcessStatus status,
Instant startedAt) { Instant startedAt) {
return new ProcessInstance(id, definitionId, initiator, status, startedAt, null, ProcessContext.empty()); return new ProcessInstance(id, definitionId, 1, initiator, status, startedAt, null, ProcessContext.empty());
} }
} }
@@ -15,7 +15,7 @@ public final class LeaveRequestExample {
public static void main(String[] args) { public static void main(String[] args) {
OrdoEngine ordo = new InMemoryOrdoEngine(); OrdoEngine ordo = new InMemoryOrdoEngine();
ordo.register(ProcessDefinition.linear("leave-request", "Leave request", List.of( ordo.publish(ProcessDefinition.linear("leave-request", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry") ApprovalStep.single("hr", "HR approval", "henry")
))); )));
@@ -3,7 +3,6 @@ package com.jetlumen.ordo.spring;
import com.jetlumen.ordo.api.OrdoEngine; import com.jetlumen.ordo.api.OrdoEngine;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessDefinitionParser; import com.jetlumen.ordo.api.ProcessDefinitionParser;
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
import org.springframework.boot.ApplicationArguments; import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner; import org.springframework.boot.ApplicationRunner;
import org.springframework.core.io.Resource; import org.springframework.core.io.Resource;
@@ -13,7 +12,7 @@ import org.springframework.core.io.support.ResourcePatternResolver;
import java.io.InputStream; import java.io.InputStream;
import java.util.Objects; import java.util.Objects;
/** Loads structural JSON process definitions and upserts them when no instance is running. */ /** Loads structural JSON process definitions and publishes a new version when the graph changed. */
public final class OrdoDefinitionLoader implements ApplicationRunner { public final class OrdoDefinitionLoader implements ApplicationRunner {
private final OrdoEngine engine; private final OrdoEngine engine;
private final String location; private final String location;
@@ -49,11 +48,7 @@ public final class OrdoDefinitionLoader implements ApplicationRunner {
} catch (RuntimeException e) { } catch (RuntimeException e) {
throw new IllegalStateException("failed to parse ordo definition from " + describe(resource), e); throw new IllegalStateException("failed to parse ordo definition from " + describe(resource), e);
} }
try { engine.publish(definition);
engine.replace(definition);
} catch (DefinitionInUseException ignored) {
// keep the stored graph while RUNNING instances still reference it
}
} }
} }
@@ -49,7 +49,7 @@ class OrdoJdbcAutoConfigurationTest {
assertThat(context).hasSingleBean(OrdoEngine.class); assertThat(context).hasSingleBean(OrdoEngine.class);
OrdoEngine engine = context.getBean(OrdoEngine.class); OrdoEngine engine = context.getBean(OrdoEngine.class);
engine.register(LEAVE_REQUEST); engine.publish(LEAVE_REQUEST);
ProcessInstance instance = engine.start("leave-request", "alice"); ProcessInstance instance = engine.start("leave-request", "alice");
List<ApprovalTask> pending = engine.findPendingTasksByAssignee("maria"); List<ApprovalTask> pending = engine.findPendingTasksByAssignee("maria");
assertThat(pending).hasSize(1); assertThat(pending).hasSize(1);
@@ -92,8 +92,9 @@ class OrdoJdbcAutoConfigurationTest {
loader.load(); loader.load();
ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class) ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class)
.findById("leave-request-routed") .findLatest("leave-request-routed")
.orElseThrow(); .orElseThrow();
assertThat(definition.version()).isEqualTo(1);
assertThat(definition.name()).isEqualTo("Leave request"); assertThat(definition.name()).isEqualTo("Leave request");
assertThat(definition.steps()).hasSize(2); assertThat(definition.steps()).hasSize(2);
assertThat(definition.transitions()).hasSize(3); assertThat(definition.transitions()).hasSize(3);
@@ -104,13 +105,13 @@ class OrdoJdbcAutoConfigurationTest {
void reloadsClasspathJsonDefinitionsWhenNoInstanceIsRunning() { void reloadsClasspathJsonDefinitionsWhenNoInstanceIsRunning() {
withDataSourceRunner.run(context -> { withDataSourceRunner.run(context -> {
OrdoEngine engine = context.getBean(OrdoEngine.class); OrdoEngine engine = context.getBean(OrdoEngine.class);
engine.replace(ProcessDefinition.linear("leave-request-routed", "stale", List.of( engine.publish(ProcessDefinition.linear("leave-request-routed", "stale", List.of(
ApprovalStep.single("lead", "Lead approval", "lee")))); ApprovalStep.single("lead", "Lead approval", "lee"))));
context.getBean(OrdoDefinitionLoader.class).load(); context.getBean(OrdoDefinitionLoader.class).load();
ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class) ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class)
.findById("leave-request-routed") .findLatest("leave-request-routed")
.orElseThrow(); .orElseThrow();
assertThat(definition.name()).isEqualTo("Leave request"); assertThat(definition.name()).isEqualTo("Leave request");
assertThat(definition.steps()).hasSize(2); assertThat(definition.steps()).hasSize(2);
@@ -118,20 +119,23 @@ class OrdoJdbcAutoConfigurationTest {
} }
@Test @Test
void keepsStoredDefinitionWhenReloadFindsRunningInstances() { void publishesClasspathJsonEvenWhenInstancesAreRunning() {
withDataSourceRunner.run(context -> { withDataSourceRunner.run(context -> {
OrdoEngine engine = context.getBean(OrdoEngine.class); OrdoEngine engine = context.getBean(OrdoEngine.class);
engine.replace(ProcessDefinition.linear("leave-request-routed", "stale", List.of( engine.publish(ProcessDefinition.linear("leave-request-routed", "stale", List.of(
ApprovalStep.single("lead", "Lead approval", "lee")))); ApprovalStep.single("lead", "Lead approval", "lee"))));
engine.start("leave-request-routed", "alice"); ProcessInstance instance = engine.start("leave-request-routed", "alice");
context.getBean(OrdoDefinitionLoader.class).load(); context.getBean(OrdoDefinitionLoader.class).load();
ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class) ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class)
.findById("leave-request-routed") .findLatest("leave-request-routed")
.orElseThrow(); .orElseThrow();
assertThat(definition.name()).isEqualTo("stale"); assertThat(definition.name()).isEqualTo("Leave request");
assertThat(definition.steps()).hasSize(1); assertThat(definition.steps()).hasSize(2);
assertThat(instance.definitionVersion()).isEqualTo(engine.findInstance(instance.id()).orElseThrow()
.definitionVersion());
assertThat(engine.findPendingTasksByInstanceId(instance.id()).get(0).stepId()).isEqualTo("lead");
}); });
} }
@@ -160,7 +164,7 @@ class OrdoJdbcAutoConfigurationTest {
withDataSourceRunner.withUserConfiguration(RecordingListenerConfig.class) withDataSourceRunner.withUserConfiguration(RecordingListenerConfig.class)
.run(context -> { .run(context -> {
OrdoEngine engine = context.getBean(OrdoEngine.class); OrdoEngine engine = context.getBean(OrdoEngine.class);
engine.register(LEAVE_REQUEST); engine.publish(LEAVE_REQUEST);
ProcessInstance instance = engine.start("leave-request", "alice"); ProcessInstance instance = engine.start("leave-request", "alice");
engine.approve(engine.findPendingTasksByAssignee("maria").get(0).id(), "maria"); engine.approve(engine.findPendingTasksByAssignee("maria").get(0).id(), "maria");
@@ -15,6 +15,8 @@ import java.sql.Connection;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.ResultSet; import java.sql.ResultSet;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.Timestamp;
import java.time.Instant;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.List; import java.util.List;
@@ -24,36 +26,45 @@ import java.util.Optional;
/** JDBC implementation of the definition storage port; steps and their candidates live in separate tables. */ /** JDBC implementation of the definition storage port; steps and their candidates live in separate tables. */
public final class JdbcProcessDefinitionRepository implements ProcessDefinitionRepository { public final class JdbcProcessDefinitionRepository implements ProcessDefinitionRepository {
private static final String INSERT_PROCESS =
"INSERT INTO ordo_process (id, current_version, name) VALUES (?, ?, ?)";
private static final String UPDATE_PROCESS =
"UPDATE ordo_process SET current_version = ?, name = ? WHERE id = ?";
private static final String LOCK_PROCESS =
"SELECT current_version FROM ordo_process WHERE id = ? FOR UPDATE";
private static final String INSERT_DEFINITION = private static final String INSERT_DEFINITION =
"INSERT INTO ordo_process_definition (id, name) VALUES (?, ?)"; "INSERT INTO ordo_process_definition (id, version, name, created_at) VALUES (?, ?, ?, ?)";
private static final String INSERT_STEP = private static final String INSERT_STEP =
"INSERT INTO ordo_approval_step (definition_id, step_id, step_name, policy, step_order, kind, action_key," "INSERT INTO ordo_approval_step (definition_id, definition_version, step_id, step_name, policy, step_order,"
+ " due_after, due_then, due_to, due_action) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"; + " kind, action_key, due_after, due_then, due_to, due_action)"
+ " VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
private static final String INSERT_CANDIDATE = private static final String INSERT_CANDIDATE =
"INSERT INTO ordo_step_candidate (definition_id, step_id, candidate, candidate_order) VALUES (?, ?, ?, ?)"; "INSERT INTO ordo_step_candidate (definition_id, definition_version, step_id, candidate, candidate_order)"
+ " VALUES (?, ?, ?, ?, ?)";
private static final String INSERT_TRANSITION = private static final String INSERT_TRANSITION =
"INSERT INTO ordo_step_transition (definition_id, from_step_id, to_step_id, condition_key, priority) VALUES (?, ?, ?, ?, ?)"; "INSERT INTO ordo_step_transition (definition_id, definition_version, from_step_id, to_step_id,"
private static final String UPDATE_DEFINITION_NAME = + " condition_key, priority) VALUES (?, ?, ?, ?, ?, ?)";
"UPDATE ordo_process_definition SET name = ? WHERE id = ?";
private static final String DELETE_TRANSITIONS =
"DELETE FROM ordo_step_transition WHERE definition_id = ?";
private static final String DELETE_CANDIDATES =
"DELETE FROM ordo_step_candidate WHERE definition_id = ?";
private static final String DELETE_STEPS =
"DELETE FROM ordo_approval_step WHERE definition_id = ?";
private static final String SELECT_DEFINITION = private static final String SELECT_DEFINITION =
"SELECT id, name FROM ordo_process_definition WHERE id = ?"; "SELECT name FROM ordo_process_definition WHERE id = ? AND version = ?";
private static final String SELECT_STEPS = private static final String SELECT_STEPS =
"SELECT step_id, step_name, policy, kind, action_key, due_after, due_then, due_to, due_action" "SELECT step_id, step_name, policy, kind, action_key, due_after, due_then, due_to, due_action"
+ " FROM ordo_approval_step" + " FROM ordo_approval_step"
+ " WHERE definition_id = ? ORDER BY step_order"; + " WHERE definition_id = ? AND definition_version = ? ORDER BY step_order";
private static final String SELECT_CANDIDATES = private static final String SELECT_CANDIDATES =
"SELECT step_id, candidate FROM ordo_step_candidate WHERE definition_id = ? ORDER BY step_id, candidate_order"; "SELECT step_id, candidate FROM ordo_step_candidate"
+ " WHERE definition_id = ? AND definition_version = ? ORDER BY step_id, candidate_order";
private static final String SELECT_TRANSITIONS = private static final String SELECT_TRANSITIONS =
"SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition " "SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition "
+ "WHERE definition_id = ? ORDER BY from_step_id, priority"; + "WHERE definition_id = ? AND definition_version = ? ORDER BY from_step_id, priority";
private static final String SELECT_DEFINITIONS_PAGE = private static final String SELECT_LATEST_PAGE =
"SELECT id, name FROM ordo_process_definition ORDER BY id LIMIT ? OFFSET ?"; "SELECT p.id, p.current_version, d.name FROM ordo_process p"
+ " JOIN ordo_process_definition d ON d.id = p.id AND d.version = p.current_version"
+ " ORDER BY p.id LIMIT ? OFFSET ?";
private static final String COUNT_PROCESSES = "SELECT COUNT(*) FROM ordo_process";
private static final String SELECT_VERSIONS_PAGE =
"SELECT version, name FROM ordo_process_definition WHERE id = ? ORDER BY version DESC LIMIT ? OFFSET ?";
private static final String COUNT_VERSIONS =
"SELECT COUNT(*) FROM ordo_process_definition WHERE id = ?";
private final JdbcConnectionProvider connectionProvider; private final JdbcConnectionProvider connectionProvider;
@@ -62,62 +73,51 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
@Override @Override
public boolean insertIfAbsent(ProcessDefinition definition) { public ProcessDefinition publish(ProcessDefinition definition) {
Objects.requireNonNull(definition, "definition must not be null"); Objects.requireNonNull(definition, "definition must not be null");
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try { try {
if (!insertDefinitionRow(connection, definition)) { Integer current = lockCurrentVersion(connection, definition.id());
return false; if (current == null) {
if (insertProcessRow(connection, definition.id(), 1, definition.name())) {
ProcessDefinition first = definition.withVersion(1);
insertDefinitionVersion(connection, first);
insertGraph(connection, first);
return first;
}
current = lockCurrentVersion(connection, definition.id());
if (current == null) {
throw new JdbcStorageException(
"failed to lock process after concurrent publish: " + definition.id(),
new IllegalStateException("process row missing"));
}
} }
insertGraph(connection, definition); ProcessDefinition latest = assemble(connection, definition.id(), current);
return true; if (latest.sameGraph(definition)) {
return latest;
}
int next = current + 1;
ProcessDefinition published = definition.withVersion(next);
insertDefinitionVersion(connection, published);
insertGraph(connection, published);
updateProcess(connection, definition.id(), next, definition.name());
return published;
} catch (SQLException e) { } catch (SQLException e) {
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e); throw new JdbcStorageException("failed to publish definition: " + definition.id(), e);
} finally { } finally {
connectionProvider.close(connection); connectionProvider.close(connection);
} }
} }
@Override @Override
public void upsert(ProcessDefinition definition) { public Optional<ProcessDefinition> findLatest(String definitionId) {
Objects.requireNonNull(definition, "definition must not be null");
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try { try {
// PostgreSQL aborts the current transaction on unique-constraint violations, Integer version = currentVersion(connection, definitionId);
// so upsert must not probe existence via a failing INSERT. if (version == null) {
if (definitionExists(connection, definition.id())) { return Optional.empty();
updateDefinitionName(connection, definition);
deleteGraph(connection, definition.id());
} else {
try (PreparedStatement insert = connection.prepareStatement(INSERT_DEFINITION)) {
insert.setString(1, definition.id());
insert.setString(2, definition.name());
insert.executeUpdate();
}
} }
insertGraph(connection, definition); return Optional.of(assemble(connection, definitionId, version));
} catch (SQLException e) {
throw new JdbcStorageException("failed to upsert definition: " + definition.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public Optional<ProcessDefinition> findById(String definitionId) {
Connection connection = connectionProvider.getConnection();
try {
String name;
try (PreparedStatement selectDefinition = connection.prepareStatement(SELECT_DEFINITION)) {
selectDefinition.setString(1, definitionId);
try (ResultSet resultSet = selectDefinition.executeQuery()) {
if (!resultSet.next()) {
return Optional.empty();
}
name = resultSet.getString("name");
}
}
return Optional.of(assemble(connection, definitionId, name));
} catch (SQLException e) { } catch (SQLException e) {
throw new JdbcStorageException("failed to load definition: " + definitionId, e); throw new JdbcStorageException("failed to load definition: " + definitionId, e);
} finally { } finally {
@@ -125,24 +125,41 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
} }
@Override
public Optional<ProcessDefinition> find(String definitionId, int version) {
Connection connection = connectionProvider.getConnection();
try {
if (!definitionExists(connection, definitionId, version)) {
return Optional.empty();
}
return Optional.of(assemble(connection, definitionId, version));
} catch (SQLException e) {
throw new JdbcStorageException("failed to load definition: " + definitionId + " version " + version, e);
} finally {
connectionProvider.close(connection);
}
}
@Override @Override
public Page<ProcessDefinition> findAll(PageRequest pageRequest) { public Page<ProcessDefinition> findAll(PageRequest pageRequest) {
Objects.requireNonNull(pageRequest, "pageRequest must not be null"); Objects.requireNonNull(pageRequest, "pageRequest must not be null");
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try { try {
long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_process_definition", List.of()); long total = PageSupport.count(connection, COUNT_PROCESSES, List.of());
List<ProcessDefinition> content = new ArrayList<>(); List<ProcessDefinition> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITIONS_PAGE)) { try (PreparedStatement select = connection.prepareStatement(SELECT_LATEST_PAGE)) {
select.setInt(1, pageRequest.size()); select.setInt(1, pageRequest.size());
select.setInt(2, pageRequest.offset()); select.setInt(2, pageRequest.offset());
List<Map.Entry<String, String>> idsAndNames = new ArrayList<>(); List<int[]> versions = new ArrayList<>();
List<String> ids = new ArrayList<>();
try (ResultSet resultSet = select.executeQuery()) { try (ResultSet resultSet = select.executeQuery()) {
while (resultSet.next()) { while (resultSet.next()) {
idsAndNames.add(Map.entry(resultSet.getString("id"), resultSet.getString("name"))); ids.add(resultSet.getString("id"));
versions.add(new int[] {resultSet.getInt("current_version")});
} }
} }
for (Map.Entry<String, String> idAndName : idsAndNames) { for (int i = 0; i < ids.size(); i++) {
content.add(assemble(connection, idAndName.getKey(), idAndName.getValue())); content.add(assemble(connection, ids.get(i), versions.get(i)[0]));
} }
} }
return new Page<>(content, total, pageRequest.page(), pageRequest.size()); return new Page<>(content, total, pageRequest.page(), pageRequest.size());
@@ -153,10 +170,52 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
} }
private static ProcessDefinition assemble(Connection connection, String definitionId, String name) throws SQLException { @Override
public Page<ProcessDefinition> findVersions(String definitionId, PageRequest pageRequest) {
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
Connection connection = connectionProvider.getConnection();
try {
long total = PageSupport.count(connection, COUNT_VERSIONS, List.of(definitionId));
List<Integer> versionNumbers = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(SELECT_VERSIONS_PAGE)) {
select.setString(1, definitionId);
select.setInt(2, pageRequest.size());
select.setInt(3, pageRequest.offset());
try (ResultSet resultSet = select.executeQuery()) {
while (resultSet.next()) {
versionNumbers.add(resultSet.getInt("version"));
}
}
}
List<ProcessDefinition> content = new ArrayList<>();
for (int version : versionNumbers) {
content.add(assemble(connection, definitionId, version));
}
return new Page<>(content, total, pageRequest.page(), pageRequest.size());
} catch (SQLException e) {
throw new JdbcStorageException("failed to list definition versions: " + definitionId, e);
} finally {
connectionProvider.close(connection);
}
}
private static ProcessDefinition assemble(Connection connection, String definitionId, int version)
throws SQLException {
String name;
try (PreparedStatement selectDefinition = connection.prepareStatement(SELECT_DEFINITION)) {
selectDefinition.setString(1, definitionId);
selectDefinition.setInt(2, version);
try (ResultSet resultSet = selectDefinition.executeQuery()) {
if (!resultSet.next()) {
throw new SQLException("definition row missing: " + definitionId + " version " + version);
}
name = resultSet.getString("name");
}
}
List<StepRow> stepRows = new ArrayList<>(); List<StepRow> stepRows = new ArrayList<>();
try (PreparedStatement selectSteps = connection.prepareStatement(SELECT_STEPS)) { try (PreparedStatement selectSteps = connection.prepareStatement(SELECT_STEPS)) {
selectSteps.setString(1, definitionId); selectSteps.setString(1, definitionId);
selectSteps.setInt(2, version);
try (ResultSet resultSet = selectSteps.executeQuery()) { try (ResultSet resultSet = selectSteps.executeQuery()) {
while (resultSet.next()) { while (resultSet.next()) {
stepRows.add(ApprovalStepMapper.readRow(resultSet)); stepRows.add(ApprovalStepMapper.readRow(resultSet));
@@ -166,6 +225,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
Map<String, List<String>> candidatesByStep = new LinkedHashMap<>(); Map<String, List<String>> candidatesByStep = new LinkedHashMap<>();
try (PreparedStatement selectCandidates = connection.prepareStatement(SELECT_CANDIDATES)) { try (PreparedStatement selectCandidates = connection.prepareStatement(SELECT_CANDIDATES)) {
selectCandidates.setString(1, definitionId); selectCandidates.setString(1, definitionId);
selectCandidates.setInt(2, version);
try (ResultSet resultSet = selectCandidates.executeQuery()) { try (ResultSet resultSet = selectCandidates.executeQuery()) {
while (resultSet.next()) { while (resultSet.next()) {
candidatesByStep.computeIfAbsent(resultSet.getString("step_id"), key -> new ArrayList<>()) candidatesByStep.computeIfAbsent(resultSet.getString("step_id"), key -> new ArrayList<>())
@@ -182,6 +242,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
List<StepTransition> transitions = new ArrayList<>(); List<StepTransition> transitions = new ArrayList<>();
try (PreparedStatement selectTransitions = connection.prepareStatement(SELECT_TRANSITIONS)) { try (PreparedStatement selectTransitions = connection.prepareStatement(SELECT_TRANSITIONS)) {
selectTransitions.setString(1, definitionId); selectTransitions.setString(1, definitionId);
selectTransitions.setInt(2, version);
try (ResultSet resultSet = selectTransitions.executeQuery()) { try (ResultSet resultSet = selectTransitions.executeQuery()) {
while (resultSet.next()) { while (resultSet.next()) {
TransitionRow row = StepTransitionMapper.readRow(resultSet); TransitionRow row = StepTransitionMapper.readRow(resultSet);
@@ -190,23 +251,46 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
} }
} }
return new ProcessDefinition(definitionId, name, steps, transitions); return new ProcessDefinition(definitionId, version, name, steps, transitions);
} }
private static boolean definitionExists(Connection connection, String definitionId) throws SQLException { private static Integer lockCurrentVersion(Connection connection, String definitionId) throws SQLException {
try (PreparedStatement select = connection.prepareStatement(LOCK_PROCESS)) {
select.setString(1, definitionId);
try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next() ? resultSet.getInt("current_version") : null;
}
}
}
private static Integer currentVersion(Connection connection, String definitionId) throws SQLException {
try (PreparedStatement select = connection.prepareStatement(
"SELECT current_version FROM ordo_process WHERE id = ?")) {
select.setString(1, definitionId);
try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next() ? resultSet.getInt("current_version") : null;
}
}
}
private static boolean definitionExists(Connection connection, String definitionId, int version)
throws SQLException {
try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITION)) { try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITION)) {
select.setString(1, definitionId); select.setString(1, definitionId);
select.setInt(2, version);
try (ResultSet resultSet = select.executeQuery()) { try (ResultSet resultSet = select.executeQuery()) {
return resultSet.next(); return resultSet.next();
} }
} }
} }
private static boolean insertDefinitionRow(Connection connection, ProcessDefinition definition) throws SQLException { private static boolean insertProcessRow(Connection connection, String id, int version, String name)
try (PreparedStatement insertDefinition = connection.prepareStatement(INSERT_DEFINITION)) { throws SQLException {
insertDefinition.setString(1, definition.id()); try (PreparedStatement insert = connection.prepareStatement(INSERT_PROCESS)) {
insertDefinition.setString(2, definition.name()); insert.setString(1, id);
insertDefinition.executeUpdate(); insert.setInt(2, version);
insert.setString(3, name);
insert.executeUpdate();
return true; return true;
} catch (SQLException e) { } catch (SQLException e) {
if (isDuplicateKey(e)) { if (isDuplicateKey(e)) {
@@ -216,26 +300,23 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
} }
private static void updateDefinitionName(Connection connection, ProcessDefinition definition) throws SQLException { private static void updateProcess(Connection connection, String id, int version, String name) throws SQLException {
try (PreparedStatement update = connection.prepareStatement(UPDATE_DEFINITION_NAME)) { try (PreparedStatement update = connection.prepareStatement(UPDATE_PROCESS)) {
update.setString(1, definition.name()); update.setInt(1, version);
update.setString(2, definition.id()); update.setString(2, name);
update.setString(3, id);
update.executeUpdate(); update.executeUpdate();
} }
} }
private static void deleteGraph(Connection connection, String definitionId) throws SQLException { private static void insertDefinitionVersion(Connection connection, ProcessDefinition definition)
try (PreparedStatement deleteTransitions = connection.prepareStatement(DELETE_TRANSITIONS)) { throws SQLException {
deleteTransitions.setString(1, definitionId); try (PreparedStatement insert = connection.prepareStatement(INSERT_DEFINITION)) {
deleteTransitions.executeUpdate(); insert.setString(1, definition.id());
} insert.setInt(2, definition.version());
try (PreparedStatement deleteCandidates = connection.prepareStatement(DELETE_CANDIDATES)) { insert.setString(3, definition.name());
deleteCandidates.setString(1, definitionId); insert.setTimestamp(4, Timestamp.from(Instant.now()));
deleteCandidates.executeUpdate(); insert.executeUpdate();
}
try (PreparedStatement deleteSteps = connection.prepareStatement(DELETE_STEPS)) {
deleteSteps.setString(1, definitionId);
deleteSteps.executeUpdate();
} }
} }
@@ -244,22 +325,23 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
for (ApprovalStep step : definition.steps()) { for (ApprovalStep step : definition.steps()) {
try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) { try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) {
insertStep.setString(1, definition.id()); insertStep.setString(1, definition.id());
insertStep.setString(2, step.id()); insertStep.setInt(2, definition.version());
insertStep.setString(3, step.name()); insertStep.setString(3, step.id());
insertStep.setString(4, step.policy().name()); insertStep.setString(4, step.name());
insertStep.setInt(5, stepOrder++); insertStep.setString(5, step.policy().name());
insertStep.setString(6, step.kind().name()); insertStep.setInt(6, stepOrder++);
insertStep.setString(7, step.actionKey()); insertStep.setString(7, step.kind().name());
insertStep.setString(8, step.actionKey());
if (step.due() == null) { if (step.due() == null) {
insertStep.setString(8, null);
insertStep.setString(9, null); insertStep.setString(9, null);
insertStep.setString(10, null); insertStep.setString(10, null);
insertStep.setString(11, null); insertStep.setString(11, null);
insertStep.setString(12, null);
} else { } else {
insertStep.setString(8, step.due().after().toString()); insertStep.setString(9, step.due().after().toString());
insertStep.setString(9, step.due().then().name()); insertStep.setString(10, step.due().then().name());
insertStep.setString(10, step.due().to()); insertStep.setString(11, step.due().to());
insertStep.setString(11, step.due().action()); insertStep.setString(12, step.due().action());
} }
insertStep.executeUpdate(); insertStep.executeUpdate();
} }
@@ -267,9 +349,10 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
for (String candidate : step.candidates()) { for (String candidate : step.candidates()) {
try (PreparedStatement insertCandidate = connection.prepareStatement(INSERT_CANDIDATE)) { try (PreparedStatement insertCandidate = connection.prepareStatement(INSERT_CANDIDATE)) {
insertCandidate.setString(1, definition.id()); insertCandidate.setString(1, definition.id());
insertCandidate.setString(2, step.id()); insertCandidate.setInt(2, definition.version());
insertCandidate.setString(3, candidate); insertCandidate.setString(3, step.id());
insertCandidate.setInt(4, candidateOrder++); insertCandidate.setString(4, candidate);
insertCandidate.setInt(5, candidateOrder++);
insertCandidate.executeUpdate(); insertCandidate.executeUpdate();
} }
} }
@@ -277,10 +360,11 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
for (StepTransition transition : definition.transitions()) { for (StepTransition transition : definition.transitions()) {
try (PreparedStatement insertTransition = connection.prepareStatement(INSERT_TRANSITION)) { try (PreparedStatement insertTransition = connection.prepareStatement(INSERT_TRANSITION)) {
insertTransition.setString(1, definition.id()); insertTransition.setString(1, definition.id());
insertTransition.setString(2, transition.fromStepId()); insertTransition.setInt(2, definition.version());
insertTransition.setString(3, transition.toStepId()); insertTransition.setString(3, transition.fromStepId());
insertTransition.setString(4, transition.conditionKey()); insertTransition.setString(4, transition.toStepId());
insertTransition.setInt(5, transition.priority()); insertTransition.setString(5, transition.conditionKey());
insertTransition.setInt(6, transition.priority());
insertTransition.executeUpdate(); insertTransition.executeUpdate();
} }
} }
@@ -22,14 +22,14 @@ import java.util.Optional;
/** JDBC implementation of the instance storage port. */ /** JDBC implementation of the instance storage port. */
public final class JdbcProcessInstanceRepository implements ProcessInstanceRepository { public final class JdbcProcessInstanceRepository implements ProcessInstanceRepository {
private static final String INSERT_INSTANCE = private static final String INSERT_INSTANCE =
"INSERT INTO ordo_process_instance (id, definition_id, initiator, status, context_json, started_at, finished_at)" "INSERT INTO ordo_process_instance (id, definition_id, definition_version, initiator, status, context_json,"
+ " VALUES (?, ?, ?, ?, ?, ?, ?)"; + " started_at, finished_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)";
private static final String UPDATE_INSTANCE = private static final String UPDATE_INSTANCE =
"UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ?"; "UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ?";
private static final String COMPLETE_IF_RUNNING = private static final String COMPLETE_IF_RUNNING =
"UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ? AND status = 'RUNNING'"; "UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ? AND status = 'RUNNING'";
private static final String SELECT_INSTANCE = private static final String SELECT_INSTANCE =
"SELECT id, definition_id, initiator, status, context_json, started_at, finished_at" "SELECT id, definition_id, definition_version, initiator, status, context_json, started_at, finished_at"
+ " FROM ordo_process_instance WHERE id = ?"; + " FROM ordo_process_instance WHERE id = ?";
private static final String EXISTS_RUNNING = private static final String EXISTS_RUNNING =
"SELECT 1 FROM ordo_process_instance WHERE definition_id = ? AND status = ? LIMIT 1"; "SELECT 1 FROM ordo_process_instance WHERE definition_id = ? AND status = ? LIMIT 1";
@@ -124,8 +124,8 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
try { try {
long total = PageSupport.count(connection, long total = PageSupport.count(connection,
"SELECT COUNT(*) FROM ordo_process_instance" + where.sql(), where.params()); "SELECT COUNT(*) FROM ordo_process_instance" + where.sql(), where.params());
String sql = "SELECT id, definition_id, initiator, status, context_json, started_at, finished_at" String sql = "SELECT id, definition_id, definition_version, initiator, status, context_json, started_at,"
+ " FROM ordo_process_instance" + where.sql() + " finished_at FROM ordo_process_instance" + where.sql()
+ " ORDER BY started_at DESC, id DESC LIMIT ? OFFSET ?"; + " ORDER BY started_at DESC, id DESC LIMIT ? OFFSET ?";
List<ProcessInstance> content = new ArrayList<>(); List<ProcessInstance> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(sql)) { try (PreparedStatement select = connection.prepareStatement(sql)) {
@@ -16,11 +16,12 @@ public final class ProcessInstanceMapper {
public static void bindInsert(PreparedStatement statement, ProcessInstance instance) throws SQLException { public static void bindInsert(PreparedStatement statement, ProcessInstance instance) throws SQLException {
statement.setString(1, instance.id()); statement.setString(1, instance.id());
statement.setString(2, instance.definitionId()); statement.setString(2, instance.definitionId());
statement.setString(3, instance.initiator()); statement.setInt(3, instance.definitionVersion());
statement.setString(4, instance.status().name()); statement.setString(4, instance.initiator());
statement.setString(5, ProcessContextCodec.encode(instance.context())); statement.setString(5, instance.status().name());
statement.setTimestamp(6, Timestamp.from(instance.startedAt())); statement.setString(6, ProcessContextCodec.encode(instance.context()));
statement.setTimestamp(7, instance.finishedAt() == null ? null : Timestamp.from(instance.finishedAt())); statement.setTimestamp(7, Timestamp.from(instance.startedAt()));
statement.setTimestamp(8, instance.finishedAt() == null ? null : Timestamp.from(instance.finishedAt()));
} }
public static void bindUpdate(PreparedStatement statement, ProcessInstance instance) throws SQLException { public static void bindUpdate(PreparedStatement statement, ProcessInstance instance) throws SQLException {
@@ -35,6 +36,7 @@ public final class ProcessInstanceMapper {
return new ProcessInstance( return new ProcessInstance(
resultSet.getString("id"), resultSet.getString("id"),
resultSet.getString("definition_id"), resultSet.getString("definition_id"),
resultSet.getInt("definition_version"),
resultSet.getString("initiator"), resultSet.getString("initiator"),
ProcessStatus.valueOf(resultSet.getString("status")), ProcessStatus.valueOf(resultSet.getString("status")),
startedAt.toInstant(), startedAt.toInstant(),
@@ -0,0 +1,96 @@
CREATE TABLE ordo_process (
id VARCHAR(64) PRIMARY KEY,
current_version INTEGER NOT NULL,
name VARCHAR(255) NOT NULL
);
INSERT INTO ordo_process (id, current_version, name)
SELECT id, 1, name FROM ordo_process_definition;
ALTER TABLE ordo_process_instance DROP CONSTRAINT fk_process_instance_definition;
ALTER TABLE ordo_approval_step DROP CONSTRAINT fk_approval_step_definition;
ALTER TABLE ordo_step_candidate DROP CONSTRAINT fk_step_candidate_step;
ALTER TABLE ordo_step_transition DROP CONSTRAINT fk_transition_from;
ALTER TABLE ordo_step_transition DROP CONSTRAINT fk_transition_to;
CREATE TABLE ordo_process_definition_v7 (
id VARCHAR(64) NOT NULL,
version INTEGER NOT NULL,
name VARCHAR(255) NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (id, version)
);
INSERT INTO ordo_process_definition_v7 (id, version, name)
SELECT id, 1, name FROM ordo_process_definition;
DROP TABLE ordo_process_definition;
ALTER TABLE ordo_process_definition_v7 RENAME TO ordo_process_definition;
CREATE TABLE ordo_approval_step_v7 (
definition_id VARCHAR(64) NOT NULL,
definition_version INTEGER NOT NULL,
step_id VARCHAR(64) NOT NULL,
step_name VARCHAR(255) NOT NULL,
step_order INTEGER NOT NULL,
policy VARCHAR(16) NOT NULL,
kind VARCHAR(16) NOT NULL,
action_key VARCHAR(255),
due_after VARCHAR(32),
due_then VARCHAR(16),
due_to VARCHAR(255),
due_action VARCHAR(255),
PRIMARY KEY (definition_id, definition_version, step_id),
CONSTRAINT fk_approval_step_definition FOREIGN KEY (definition_id, definition_version)
REFERENCES ordo_process_definition (id, version)
);
INSERT INTO ordo_approval_step_v7 (definition_id, definition_version, step_id, step_name, step_order, policy, kind,
action_key, due_after, due_then, due_to, due_action)
SELECT definition_id, 1, step_id, step_name, step_order, policy, kind, action_key, due_after, due_then, due_to,
due_action
FROM ordo_approval_step;
DROP TABLE ordo_approval_step;
ALTER TABLE ordo_approval_step_v7 RENAME TO ordo_approval_step;
CREATE TABLE ordo_step_candidate_v7 (
definition_id VARCHAR(64) NOT NULL,
definition_version INTEGER NOT NULL,
step_id VARCHAR(64) NOT NULL,
candidate VARCHAR(255) NOT NULL,
candidate_order INTEGER NOT NULL,
PRIMARY KEY (definition_id, definition_version, step_id, candidate),
CONSTRAINT fk_step_candidate_step FOREIGN KEY (definition_id, definition_version, step_id)
REFERENCES ordo_approval_step (definition_id, definition_version, step_id)
);
INSERT INTO ordo_step_candidate_v7 (definition_id, definition_version, step_id, candidate, candidate_order)
SELECT definition_id, 1, step_id, candidate, candidate_order FROM ordo_step_candidate;
DROP TABLE ordo_step_candidate;
ALTER TABLE ordo_step_candidate_v7 RENAME TO ordo_step_candidate;
CREATE TABLE ordo_step_transition_v7 (
definition_id VARCHAR(64) NOT NULL,
definition_version INTEGER NOT NULL,
from_step_id VARCHAR(64) NOT NULL,
to_step_id VARCHAR(64),
condition_key VARCHAR(255),
priority INTEGER NOT NULL,
PRIMARY KEY (definition_id, definition_version, from_step_id, priority),
CONSTRAINT fk_transition_from FOREIGN KEY (definition_id, definition_version, from_step_id)
REFERENCES ordo_approval_step (definition_id, definition_version, step_id),
CONSTRAINT fk_transition_to FOREIGN KEY (definition_id, definition_version, to_step_id)
REFERENCES ordo_approval_step (definition_id, definition_version, step_id)
);
INSERT INTO ordo_step_transition_v7 (definition_id, definition_version, from_step_id, to_step_id, condition_key, priority)
SELECT definition_id, 1, from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition;
DROP TABLE ordo_step_transition;
ALTER TABLE ordo_step_transition_v7 RENAME TO ordo_step_transition;
ALTER TABLE ordo_process_instance ADD COLUMN definition_version INTEGER DEFAULT 1 NOT NULL;
ALTER TABLE ordo_process_instance ADD CONSTRAINT fk_process_instance_definition
FOREIGN KEY (definition_id, definition_version) REFERENCES ordo_process_definition (id, version);
@@ -28,9 +28,9 @@ class JdbcActionExecutionRepositoryTest {
void setUp() { void setUp() {
JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
repository = new JdbcActionExecutionRepository(connectionProvider); repository = new JdbcActionExecutionRepository(connectionProvider);
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("leave",
"Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", "alice", new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", 1, "alice",
ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); ProcessStatus.RUNNING, T0, null, ProcessContext.empty()));
} }
@@ -42,12 +42,12 @@ class JdbcApprovalTaskRepositoryTest {
/** Tasks reference their instance, which references its definition; both parent rows must exist. */ /** Tasks reference their instance, which references its definition; both parent rows must exist. */
private void insertFixtureData() { private void insertFixtureData() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("leave",
"Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider); JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider);
instanceRepository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, instanceRepository.insert(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty())); CREATED_AT, null, ProcessContext.empty()));
instanceRepository.insert(new ProcessInstance("inst-2", "leave", "alice", ProcessStatus.RUNNING, instanceRepository.insert(new ProcessInstance("inst-2", "leave", 1, "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty())); CREATED_AT, null, ProcessContext.empty()));
} }
@@ -101,10 +101,10 @@ class JdbcApprovalTaskRepositoryTest {
@Test @Test
void queryFiltersPaginatesAndOrdersNewestFirst() { void queryFiltersPaginatesAndOrdersNewestFirst() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("expense", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("expense",
"Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank")))); "Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank"))));
JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider); JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider);
instanceRepository.insert(new ProcessInstance("inst-3", "expense", "alice", ProcessStatus.RUNNING, instanceRepository.insert(new ProcessInstance("inst-3", "expense", 1, "alice", ProcessStatus.RUNNING,
CREATED_AT, null, ProcessContext.empty())); CREATED_AT, null, ProcessContext.empty()));
ApprovalTask t1 = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT); ApprovalTask t1 = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
@@ -19,7 +19,6 @@ import com.jetlumen.ordo.api.StepKind;
import com.jetlumen.ordo.api.StepDue; import com.jetlumen.ordo.api.StepDue;
import com.jetlumen.ordo.api.StepTransition; import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.TaskStatus; import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException; import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
import com.jetlumen.ordo.api.exception.NoRouteFoundException; import com.jetlumen.ordo.api.exception.NoRouteFoundException;
@@ -49,7 +48,7 @@ class JdbcOrdoEngineIntegrationTest {
void setUp() { void setUp() {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
engine = newEngine(AssigneeResolver.direct()); engine = newEngine(AssigneeResolver.direct());
engine.register(ProcessDefinition.linear("leave", "Leave request", List.of( engine.publish(ProcessDefinition.linear("leave", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")))); ApprovalStep.single("hr", "HR approval", "henry"))));
} }
@@ -136,7 +135,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcProcessHistoryRepository(connectionProvider), new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider), new JdbcActionExecutionRepository(connectionProvider),
List.of()); List.of());
startEngine.register(dueDefinition); startEngine.publish(dueDefinition);
ProcessInstance instance = startEngine.start("leave-due", "alice"); ProcessInstance instance = startEngine.start("leave-due", "alice");
assertEquals(0, startEngine.processDue(10)); assertEquals(0, startEngine.processDue(10));
@@ -165,7 +164,7 @@ class JdbcOrdoEngineIntegrationTest {
} }
return candidate; return candidate;
}); });
failingEngine.register(ProcessDefinition.linear("leave2", "Leave request", List.of( failingEngine.publish(ProcessDefinition.linear("leave2", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")))); ApprovalStep.single("hr", "HR approval", "henry"))));
@@ -192,7 +191,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcProcessHistoryRepository(connectionProvider), new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider), new JdbcActionExecutionRepository(connectionProvider),
List.of()); List.of());
failingEngine.register(new ProcessDefinition("leave-noroute", "Leave request", failingEngine.publish(new ProcessDefinition("leave-noroute", "Leave request",
List.of(ApprovalStep.single("manager", "Manager approval", "maria")), List.of(ApprovalStep.single("manager", "Manager approval", "maria")),
List.of(StepTransition.endWhen("manager", "never", 0)))); List.of(StepTransition.endWhen("manager", "never", 0))));
@@ -223,19 +222,15 @@ class JdbcOrdoEngineIntegrationTest {
} }
@Test @Test
void replaceSwapsTheGraphWhenNoInstanceIsRunning() { void publishCreatesANewVersionWhileRunningInstancesKeepTheOldGraph() {
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of( ProcessInstance running = engine.start("leave", "alice");
engine.publish(ProcessDefinition.linear("leave", "Leave request v2", List.of(
ApprovalStep.single("director", "Director approval", "diana")))); ApprovalStep.single("director", "Director approval", "diana"))));
ProcessInstance instance = engine.start("leave", "alice"); assertEquals("maria", engine.findPendingTasksByInstanceId(running.id()).get(0).assignee());
assertEquals("diana", engine.findPendingTasksByInstanceId(instance.id()).get(0).assignee()); ProcessInstance next = engine.start("leave", "bob");
} assertEquals(2, next.definitionVersion());
assertEquals("diana", engine.findPendingTasksByInstanceId(next.id()).get(0).assignee());
@Test
void replaceIsRejectedWhileAnInstanceIsRunning() {
engine.start("leave", "alice");
assertThrows(DefinitionInUseException.class, () -> engine.replace(ProcessDefinition.linear(
"leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana")))));
} }
@Test @Test
@@ -304,7 +299,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcProcessHistoryRepository(connectionProvider), new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider), new JdbcActionExecutionRepository(connectionProvider),
List.of()); List.of());
actionEngine.register(new ProcessDefinition("leave-action", "Leave request", List.of( actionEngine.publish(new ProcessDefinition("leave-action", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.action("notify", "Notify", "ok-mail"), ApprovalStep.action("notify", "Notify", "ok-mail"),
ApprovalStep.action("fail", "Fail", "boom")), ApprovalStep.action("fail", "Fail", "boom")),
@@ -104,7 +104,7 @@ class JdbcPostgresIntegrationTest {
@Test @Test
void engineCompletesASequentialApprovalProcessOverPostgres() { void engineCompletesASequentialApprovalProcessOverPostgres() {
OrdoEngine engine = newEngine(AssigneeResolver.direct()); OrdoEngine engine = newEngine(AssigneeResolver.direct());
engine.register(ProcessDefinition.linear("leave-pg", "Leave request", List.of( engine.publish(ProcessDefinition.linear("leave-pg", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")))); ApprovalStep.single("hr", "HR approval", "henry"))));
@@ -122,23 +122,25 @@ class JdbcPostgresIntegrationTest {
} }
@Test @Test
void rejectsDuplicateDefinitionIdsViaTheDatabaseUniqueConstraint() { void publishIsIdempotentForTheSameGraphAndVersionsAChangedGraph() {
JdbcProcessDefinitionRepository repository = new JdbcProcessDefinitionRepository(connectionProvider); JdbcProcessDefinitionRepository repository = new JdbcProcessDefinitionRepository(connectionProvider);
ProcessDefinition definition = ProcessDefinition.linear("leave-dup-pg", "Leave request", ProcessDefinition definition = ProcessDefinition.linear("leave-dup-pg", "Leave request",
List.of(ApprovalStep.single("manager", "Manager approval", "maria"))); List.of(ApprovalStep.single("manager", "Manager approval", "maria")));
assertTrue(repository.insertIfAbsent(definition)); assertEquals(1, repository.publish(definition).version());
assertFalse(repository.insertIfAbsent(ProcessDefinition.linear("leave-dup-pg", "Second attempt", assertEquals(1, repository.publish(definition).version());
List.of(ApprovalStep.single("manager", "Manager approval", "maria"))))); assertEquals(2, repository.publish(ProcessDefinition.linear("leave-dup-pg", "Second attempt",
assertEquals("Leave request", repository.findById("leave-dup-pg").orElseThrow().name()); List.of(ApprovalStep.single("manager", "Manager approval", "maria")))).version());
assertEquals("Second attempt", repository.findLatest("leave-dup-pg").orElseThrow().name());
assertEquals("Leave request", repository.find("leave-dup-pg", 1).orElseThrow().name());
} }
@Test @Test
void rejectsDuplicateInstanceIdsViaTheDatabaseUniqueConstraint() { void rejectsDuplicateInstanceIdsViaTheDatabaseUniqueConstraint() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear( new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear(
"leave-dup-inst-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "leave-dup-inst-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider); JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider);
ProcessInstance instance = new ProcessInstance("inst-dup-pg", "leave-dup-inst-pg", "alice", ProcessInstance instance = new ProcessInstance("inst-dup-pg", "leave-dup-inst-pg", 1, "alice",
ProcessStatus.RUNNING, NOW, null, ProcessContext.empty()); ProcessStatus.RUNNING, NOW, null, ProcessContext.empty());
repository.insert(instance); repository.insert(instance);
@@ -147,18 +149,18 @@ class JdbcPostgresIntegrationTest {
@Test @Test
void roundsTimestampsToMicrosecondPrecision() { void roundsTimestampsToMicrosecondPrecision() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear( new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear(
"leave-time-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "leave-time-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider); JdbcProcessInstanceRepository repository = new JdbcProcessInstanceRepository(connectionProvider);
Instant microAligned = Instant.parse("2026-01-15T09:00:00.123456Z"); Instant microAligned = Instant.parse("2026-01-15T09:00:00.123456Z");
repository.insert(new ProcessInstance("inst-time-1", "leave-time-pg", "alice", repository.insert(new ProcessInstance("inst-time-1", "leave-time-pg", 1, "alice",
ProcessStatus.RUNNING, microAligned, null, ProcessContext.empty())); ProcessStatus.RUNNING, microAligned, null, ProcessContext.empty()));
assertEquals(microAligned, repository.findById("inst-time-1").orElseThrow().startedAt()); assertEquals(microAligned, repository.findById("inst-time-1").orElseThrow().startedAt());
// PostgreSQL's TIMESTAMP stores microseconds and rounds the fractional seconds // PostgreSQL's TIMESTAMP stores microseconds and rounds the fractional seconds
Instant withNanos = Instant.parse("2026-01-15T09:00:00.123456789Z"); Instant withNanos = Instant.parse("2026-01-15T09:00:00.123456789Z");
repository.insert(new ProcessInstance("inst-time-2", "leave-time-pg", "alice", repository.insert(new ProcessInstance("inst-time-2", "leave-time-pg", 1, "alice",
ProcessStatus.RUNNING, withNanos, null, ProcessContext.empty())); ProcessStatus.RUNNING, withNanos, null, ProcessContext.empty()));
assertEquals(Instant.parse("2026-01-15T09:00:00.123457Z"), assertEquals(Instant.parse("2026-01-15T09:00:00.123457Z"),
repository.findById("inst-time-2").orElseThrow().startedAt()); repository.findById("inst-time-2").orElseThrow().startedAt());
@@ -166,10 +168,10 @@ class JdbcPostgresIntegrationTest {
@Test @Test
void onlyOneOfTwoConcurrentCompletionsWinsOnPostgres() throws Exception { void onlyOneOfTwoConcurrentCompletionsWinsOnPostgres() throws Exception {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear( new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear(
"leave-race-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "leave-race-pg", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider); JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider);
instanceRepository.insert(new ProcessInstance("inst-race-pg", "leave-race-pg", "alice", instanceRepository.insert(new ProcessInstance("inst-race-pg", "leave-race-pg", 1, "alice",
ProcessStatus.RUNNING, NOW, null, ProcessContext.empty())); ProcessStatus.RUNNING, NOW, null, ProcessContext.empty()));
JdbcApprovalTaskRepository taskRepository = new JdbcApprovalTaskRepository(connectionProvider); JdbcApprovalTaskRepository taskRepository = new JdbcApprovalTaskRepository(connectionProvider);
@@ -202,7 +204,7 @@ class JdbcPostgresIntegrationTest {
} }
return candidate; return candidate;
}); });
failingEngine.register(ProcessDefinition.linear("leave-rollback-pg", "Leave request", List.of( failingEngine.publish(ProcessDefinition.linear("leave-rollback-pg", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")))); ApprovalStep.single("hr", "HR approval", "henry"))));
@@ -32,35 +32,39 @@ class JdbcProcessDefinitionRepositoryTest {
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry"))); ApprovalStep.single("hr", "HR approval", "henry")));
assertTrue(repository.insertIfAbsent(definition)); ProcessDefinition published = repository.publish(definition);
assertEquals(definition, repository.findById("leave").orElseThrow()); assertEquals(1, published.version());
assertEquals(definition.withVersion(1), repository.findLatest("leave").orElseThrow());
} }
@Test @Test
void rejectsAnExistingDefinitionId() { void publishIsIdempotentWhenTheGraphIsUnchanged() {
assertTrue(repository.insertIfAbsent(definition("leave", "Leave request v1"))); ProcessDefinition first = repository.publish(definition("leave", "Leave request v1"));
assertFalse(repository.insertIfAbsent(definition("leave", "Leave request v2"))); ProcessDefinition second = repository.publish(definition("leave", "Leave request v1"));
assertEquals(first, second);
assertEquals("Leave request v1", repository.findById("leave").orElseThrow().name()); assertEquals(1, repository.findVersions("leave", new PageRequest(0, 10)).totalElements());
} }
@Test @Test
void upsertReplacesNameStepsCandidatesAndTransitions() { void publishInsertsANewVersionWithoutRewritingTheOldGraph() {
repository.insertIfAbsent(ProcessDefinition.linear("leave", "Leave request v1", List.of( repository.publish(ProcessDefinition.linear("leave", "Leave request v1", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"), ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")))); ApprovalStep.single("hr", "HR approval", "henry"))));
ProcessDefinition replacement = new ProcessDefinition("leave", "Leave request v2", List.of( ProcessDefinition replacement = new ProcessDefinition("leave", "Leave request v2", List.of(
ApprovalStep.single("director", "Director approval", "diana")), ApprovalStep.single("director", "Director approval", "diana")),
List.of(StepTransition.end("director"))); List.of(StepTransition.end("director")));
repository.upsert(replacement); ProcessDefinition published = repository.publish(replacement);
assertEquals(replacement, repository.findById("leave").orElseThrow()); assertEquals(replacement.withVersion(2), published);
assertEquals(published, repository.findLatest("leave").orElseThrow());
assertEquals("Leave request v1", repository.find("leave", 1).orElseThrow().name());
} }
@Test @Test
void returnsEmptyForAnUnknownDefinition() { void returnsEmptyForAnUnknownDefinition() {
assertTrue(repository.findById("missing").isEmpty()); assertTrue(repository.findLatest("missing").isEmpty());
assertTrue(repository.find("missing", 1).isEmpty());
} }
@Test @Test
@@ -73,8 +77,8 @@ class JdbcProcessDefinitionRepositoryTest {
new StepTransition("manager", null, null, 1), new StepTransition("manager", null, null, 1),
StepTransition.end("director"))); StepTransition.end("director")));
assertTrue(repository.insertIfAbsent(definition)); assertEquals(definition.withVersion(1), repository.publish(definition));
assertEquals(definition, repository.findById("expense").orElseThrow()); assertEquals(definition.withVersion(1), repository.findLatest("expense").orElseThrow());
} }
@Test @Test
@@ -86,10 +90,10 @@ class JdbcProcessDefinitionRepositoryTest {
StepTransition.always("manager", "mail"), StepTransition.always("manager", "mail"),
StepTransition.end("mail"))); StepTransition.end("mail")));
assertTrue(repository.insertIfAbsent(definition)); repository.publish(definition);
assertEquals(definition, repository.findById("notify").orElseThrow()); assertEquals(definition.withVersion(1), repository.findLatest("notify").orElseThrow());
assertEquals(StepKind.ACTION, repository.findById("notify").orElseThrow().steps().get(1).kind()); assertEquals(StepKind.ACTION, repository.findLatest("notify").orElseThrow().steps().get(1).kind());
assertEquals("leave-approved-mail", repository.findById("notify").orElseThrow().steps().get(1).actionKey()); assertEquals("leave-approved-mail", repository.findLatest("notify").orElseThrow().steps().get(1).actionKey());
} }
private static ProcessDefinition definition(String id, String name) { private static ProcessDefinition definition(String id, String name) {
@@ -97,28 +101,41 @@ class JdbcProcessDefinitionRepositoryTest {
} }
@Test @Test
void findAllPaginatesDefinitionsOrderedById() { void findAllPaginatesLatestDefinitionsOrderedById() {
repository.insertIfAbsent(definition("c-def", "C")); repository.publish(definition("c-def", "C"));
repository.insertIfAbsent(definition("a-def", "A")); repository.publish(definition("a-def", "A"));
repository.insertIfAbsent(definition("b-def", "B")); repository.publish(definition("b-def", "B"));
repository.publish(definition("a-def", "A2"));
Page<ProcessDefinition> pageOne = repository.findAll(new PageRequest(0, 2)); Page<ProcessDefinition> pageOne = repository.findAll(new PageRequest(0, 2));
assertEquals(3, pageOne.totalElements()); assertEquals(3, pageOne.totalElements());
assertEquals(2, pageOne.totalPages()); assertEquals(2, pageOne.totalPages());
assertTrue(pageOne.hasNext()); assertTrue(pageOne.hasNext());
assertEquals(List.of("a-def", "b-def"), pageOne.content().stream().map(ProcessDefinition::id).toList()); assertEquals(List.of("a-def", "b-def"), pageOne.content().stream().map(ProcessDefinition::id).toList());
assertEquals(2, pageOne.content().get(0).version());
Page<ProcessDefinition> pageTwo = repository.findAll(new PageRequest(1, 2)); Page<ProcessDefinition> pageTwo = repository.findAll(new PageRequest(1, 2));
assertEquals(List.of("c-def"), pageTwo.content().stream().map(ProcessDefinition::id).toList()); assertEquals(List.of("c-def"), pageTwo.content().stream().map(ProcessDefinition::id).toList());
assertFalse(pageTwo.hasNext()); assertFalse(pageTwo.hasNext());
} }
@Test
void findVersionsPaginatesNewestFirst() {
repository.publish(definition("leave", "v1"));
repository.publish(definition("leave", "v2"));
repository.publish(definition("leave", "v3"));
Page<ProcessDefinition> page = repository.findVersions("leave", new PageRequest(0, 2));
assertEquals(3, page.totalElements());
assertEquals(List.of(3, 2), page.content().stream().map(ProcessDefinition::version).toList());
}
@Test @Test
void insertsAndReadsBackStepDue() { void insertsAndReadsBackStepDue() {
ProcessDefinition definition = ProcessDefinition.linear("leave-due", "Leave request", List.of( ProcessDefinition definition = ProcessDefinition.linear("leave-due", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY,
StepKind.APPROVAL, null, StepDue.reassign(java.time.Duration.parse("PT48H"), "director")))); StepKind.APPROVAL, null, StepDue.reassign(java.time.Duration.parse("PT48H"), "director"))));
assertTrue(repository.insertIfAbsent(definition)); repository.publish(definition);
assertEquals(definition, repository.findById("leave-due").orElseThrow()); assertEquals(definition.withVersion(1), repository.findLatest("leave-due").orElseThrow());
} }
} }
@@ -27,11 +27,11 @@ class JdbcProcessHistoryRepositoryTest {
void setUp() { void setUp() {
JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
repository = new JdbcProcessHistoryRepository(connectionProvider); repository = new JdbcProcessHistoryRepository(connectionProvider);
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("leave",
"Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", "alice", new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", 1, "alice",
ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); ProcessStatus.RUNNING, T0, null, ProcessContext.empty()));
new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-2", "leave", "bob", new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-2", "leave", 1, "bob",
ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); ProcessStatus.RUNNING, T0, null, ProcessContext.empty()));
} }
@@ -31,7 +31,7 @@ class JdbcProcessInstanceRepositoryTest {
connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource());
repository = new JdbcProcessInstanceRepository(connectionProvider); repository = new JdbcProcessInstanceRepository(connectionProvider);
// instances reference their definition, so the parent row must exist // instances reference their definition, so the parent row must exist
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("leave",
"Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria"))));
} }
@@ -43,7 +43,7 @@ class JdbcProcessInstanceRepositoryTest {
"urgent", true, "urgent", true,
"candidates", List.of("maria", "henry"), "candidates", List.of("maria", "henry"),
"meta", Map.of("priority", "high", "retries", 2))); "meta", Map.of("priority", "high", "retries", 2)));
ProcessInstance instance = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, ProcessInstance instance = new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT, null, context); STARTED_AT, null, context);
repository.insert(instance); repository.insert(instance);
@@ -53,11 +53,11 @@ class JdbcProcessInstanceRepositoryTest {
@Test @Test
void updatesStatusAndFinishedAt() { void updatesStatusAndFinishedAt() {
repository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, repository.insert(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty())); STARTED_AT, null, ProcessContext.empty()));
Instant finishedAt = STARTED_AT.plusSeconds(300); Instant finishedAt = STARTED_AT.plusSeconds(300);
repository.update(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.APPROVED, repository.update(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.APPROVED,
STARTED_AT, finishedAt, ProcessContext.empty())); STARTED_AT, finishedAt, ProcessContext.empty()));
ProcessInstance updated = repository.findById("inst-1").orElseThrow(); ProcessInstance updated = repository.findById("inst-1").orElseThrow();
@@ -74,51 +74,51 @@ class JdbcProcessInstanceRepositoryTest {
void existsRunningOnlyCountsRunningInstancesOfThatDefinition() { void existsRunningOnlyCountsRunningInstancesOfThatDefinition() {
assertFalse(repository.existsRunning("leave")); assertFalse(repository.existsRunning("leave"));
repository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, repository.insert(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty())); STARTED_AT, null, ProcessContext.empty()));
assertTrue(repository.existsRunning("leave")); assertTrue(repository.existsRunning("leave"));
assertFalse(repository.existsRunning("other")); assertFalse(repository.existsRunning("other"));
repository.update(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.APPROVED, repository.update(new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.APPROVED,
STARTED_AT, STARTED_AT.plusSeconds(60), ProcessContext.empty())); STARTED_AT, STARTED_AT.plusSeconds(60), ProcessContext.empty()));
assertFalse(repository.existsRunning("leave")); assertFalse(repository.existsRunning("leave"));
} }
@Test @Test
void completeIfRunningOnlyUpdatesARunningInstance() { void completeIfRunningOnlyUpdatesARunningInstance() {
ProcessInstance running = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, ProcessInstance running = new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty()); STARTED_AT, null, ProcessContext.empty());
repository.insert(running); repository.insert(running);
Instant finishedAt = STARTED_AT.plusSeconds(30); Instant finishedAt = STARTED_AT.plusSeconds(30);
ProcessInstance withdrawn = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.WITHDRAWN, ProcessInstance withdrawn = new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.WITHDRAWN,
STARTED_AT, finishedAt, ProcessContext.empty()); STARTED_AT, finishedAt, ProcessContext.empty());
assertTrue(repository.completeIfRunning(withdrawn)); assertTrue(repository.completeIfRunning(withdrawn));
assertEquals(ProcessStatus.WITHDRAWN, repository.findById("inst-1").orElseThrow().status()); assertEquals(ProcessStatus.WITHDRAWN, repository.findById("inst-1").orElseThrow().status());
assertEquals(finishedAt, repository.findById("inst-1").orElseThrow().finishedAt()); assertEquals(finishedAt, repository.findById("inst-1").orElseThrow().finishedAt());
assertFalse(repository.completeIfRunning(new ProcessInstance("inst-1", "leave", "alice", assertFalse(repository.completeIfRunning(new ProcessInstance("inst-1", "leave", 1, "alice",
ProcessStatus.APPROVED, STARTED_AT, finishedAt, ProcessContext.empty()))); ProcessStatus.APPROVED, STARTED_AT, finishedAt, ProcessContext.empty())));
assertFalse(repository.completeIfRunning(new ProcessInstance("missing", "leave", "alice", assertFalse(repository.completeIfRunning(new ProcessInstance("missing", "leave", 1, "alice",
ProcessStatus.WITHDRAWN, STARTED_AT, finishedAt, ProcessContext.empty()))); ProcessStatus.WITHDRAWN, STARTED_AT, finishedAt, ProcessContext.empty())));
} }
@Test @Test
void updateOfAnUnknownInstanceFails() { void updateOfAnUnknownInstanceFails() {
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave", assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",
"alice", ProcessStatus.APPROVED, STARTED_AT, STARTED_AT, ProcessContext.empty()))); 1, "alice", ProcessStatus.APPROVED, STARTED_AT, STARTED_AT, ProcessContext.empty())));
} }
@Test @Test
void queryFiltersPaginatesAndOrdersNewestFirst() { void queryFiltersPaginatesAndOrdersNewestFirst() {
new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("expense", new JdbcProcessDefinitionRepository(connectionProvider).publish(ProcessDefinition.linear("expense",
"Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank")))); "Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank"))));
ProcessInstance i1 = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, ProcessInstance i1 = new ProcessInstance("inst-1", "leave", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty()); STARTED_AT, null, ProcessContext.empty());
ProcessInstance i2 = new ProcessInstance("inst-2", "leave", "bob", ProcessStatus.RUNNING, ProcessInstance i2 = new ProcessInstance("inst-2", "leave", 1, "bob", ProcessStatus.RUNNING,
STARTED_AT.plusSeconds(5), null, ProcessContext.empty()); STARTED_AT.plusSeconds(5), null, ProcessContext.empty());
ProcessInstance i3 = new ProcessInstance("inst-3", "expense", "alice", ProcessStatus.RUNNING, ProcessInstance i3 = new ProcessInstance("inst-3", "expense", 1, "alice", ProcessStatus.RUNNING,
STARTED_AT.plusSeconds(10), null, ProcessContext.empty()); STARTED_AT.plusSeconds(10), null, ProcessContext.empty());
repository.insert(i1); repository.insert(i1);
repository.insert(i2); repository.insert(i2);
@@ -19,7 +19,8 @@ final class JdbcTestSupport {
"/db/migration/V3__add_step_transitions.sql", "/db/migration/V3__add_step_transitions.sql",
"/db/migration/V4__add_step_kind_and_action_key.sql", "/db/migration/V4__add_step_kind_and_action_key.sql",
"/db/migration/V5__add_process_event_and_action_execution.sql", "/db/migration/V5__add_process_event_and_action_execution.sql",
"/db/migration/V6__add_step_due.sql" "/db/migration/V6__add_step_due.sql",
"/db/migration/V7__definition_versions.sql"
}; };
private static final String[] SCHEMA_SQL = loadSchemas(); private static final String[] SCHEMA_SQL = loadSchemas();
@@ -60,7 +60,7 @@ class JdbcTransactionExecutorTest {
private void insertDefinition(String id, String name) { private void insertDefinition(String id, String name) {
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try (PreparedStatement insert = connection.prepareStatement( try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO ordo_process_definition (id, name) VALUES (?, ?)")) { "INSERT INTO ordo_process_definition (id, version, name, created_at) VALUES (?, 1, ?, CURRENT_TIMESTAMP)")) {
insert.setString(1, id); insert.setString(1, id);
insert.setString(2, name); insert.setString(2, name);
insert.executeUpdate(); insert.executeUpdate();