feat: add structured PARALLEL blocks with join-all tokens

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-18 18:17:53 +08:00
co-authored by Cursor
parent 0e9bb4bd16
commit e4478656c1
31 changed files with 1300 additions and 62 deletions
+2 -1
View File
@@ -20,6 +20,7 @@
- 线性或多步图:`StepTransition` + 内置谓词 / 宿主 `RoutingCondition`
- 会签/或签:`ApprovalPolicy.ALL` / `ANY`(多候选人)
- 结构化并行:`kind: PARALLEL` 块(一层、join ALL)
- ACTION 步骤:事务提交后调用宿主 `ActionHandler`
- 发起人撤回:`WITHDRAWN`,待办任务 `SKIPPED`
- 管理员/系统取消:`cancel` → `CANCELLED`,待办任务 `SKIPPED`(引擎不鉴权角色)
@@ -80,7 +81,7 @@ ordo.approve(manager.id(), "maria", "ok");
}
```
- `kind` 默认 `APPROVAL`;ACTION 用 `"action"` 作为 handler 查找键。
- `kind` 默认 `APPROVAL`;ACTION 用 `"action"` 作为 handler 查找键;`PARALLEL` 用 `branches`(结构化并行块,见 [docs/usage.md](docs/usage.md))。
- 带 `when` 的边先按 `priority` 匹配,都未命中再走无条件边;`to: null` 表示结束。
- 代码侧可用 `ProcessDefinitionParser.fromJson(...)`。
- `replace` 会整体替换同 id 定义;存在 `RUNNING` 实例时拒绝替换。
+34
View File
@@ -486,6 +486,40 @@ components:
type: array
items:
type: object
properties:
id:
type: string
name:
type: string
kind:
type: string
enum: [APPROVAL, ACTION, PARALLEL]
candidates:
type: array
items:
type: string
policy:
type: string
enum: [ANY, ALL]
action:
type: string
due:
type: object
branches:
type: array
items:
type: object
properties:
id:
type: string
steps:
type: array
items:
type: object
transitions:
type: array
items:
type: object
transitions:
type: array
items:
+3 -2
View File
@@ -7,6 +7,7 @@
## 已完成
- ANY/ALL 会签/或签
- 结构化 PARALLEL 块(一层、join ALL、驳回否决整单)
- JDBC 存储(PostgreSQL / MySQL)+ 按方言 Flyway 基线
- 定义不可变多版本:`publish`;实例锁定 `definitionVersion`
- 条件路由 `StepTransition`:内置谓词 AST 或宿主 `RoutingCondition`(`ref` + `args`)
@@ -25,7 +26,7 @@
### 独立设计器(不进本仓库)
流程设计器是**单独产品**,消费上述 REST 与目录,不做成 ordo 模块。画布对齐引擎图(审批步、ACTION 步、边上的 `when`/`priority`,结束为 `to: null`),不引入 BPMN 网关/并行等引擎没有的语义。节点坐标等 layout 由设计器自存,不进入 `ProcessDefinition`。
流程设计器是**单独产品**,消费上述 REST 与目录,不做成 ordo 模块。画布对齐引擎图(审批步、ACTION 步、PARALLEL 块、边上的 `when`/`priority`,结束为 `to: null`),不引入自由 fork-join 网关。节点坐标等 layout 由设计器自存,不进入 `ProcessDefinition`。
启动时机:REST 契约与目录 SPI 落地之后。草稿由设计器/REST 文档存储,不进入引擎图版本。
@@ -36,6 +37,6 @@
## 后续新特性
上表「确定要做」之外,仍可能立项其他能力(例如表单/UI schema、子流程、并行 fork-join)。**多租户在另有明确决定前不进入计划。**
上表「确定要做」之外,仍可能立项其他能力(例如表单/UI schema、子流程、嵌套 PARALLEL / join ANY)。**多租户在另有明确决定前不进入计划。**
新特性立项时写入「开发计划」对应小节;完成后移到「已完成」,并更新 [usage.md](usage.md)。
+15 -3
View File
@@ -205,7 +205,8 @@ new ProcessDefinition("leave-request-routed", "Leave request",
| 字段 | 说明 |
|---|---|
| `startStep` | 可选。指定起始步骤 id;缺省为 `steps` 数组第一项。解析时会把该步旋到列表首位。 |
| `kind` | `APPROVAL`(默认)或 `ACTION`。 |
| `kind` | `APPROVAL`(默认)、`ACTION` 或 `PARALLEL`。 |
| `branches` | 仅 `PARALLEL`。至少 2 条分支;每条含 `id`、`steps`、`transitions`。禁止嵌套 PARALLEL。 |
| `candidates` / `policy` | 仅审批步。`policy` 默认 `ANY`。审批步至少一名候选人,禁止重复。 |
| `action` | 仅 ACTION 步,对应 `ActionHandler.execute` 的 key。审批步禁止带 `action`。 |
| `due` | 仅审批步。可选。`after` 为 ISO-8601 时长;`then` 为 `reassign` / `notify` / `goto`。进入该步时任务 `dueAt = now + after`。 |
@@ -218,6 +219,7 @@ new ProcessDefinition("leave-request-routed", "Leave request",
```
start → RUNNING
审批步:为每个候选人建 PENDING 任务
PARALLEL 步:为每条分支建立令牌并同时进入分支起点
ACTION 步:事务内记 PENDING 执行记录,提交后调 ActionHandler
approve / 路由结束 → APPROVED
reject(步被否决)→ REJECTED
@@ -270,7 +272,17 @@ ordo.cancel(instance.id(), "admin", "政策变更");
候选人创建任务前会经过 `AssigneeResolver.resolve(candidate, step, context)`,例如把角色名解析成用户 id。默认实现原样返回 candidate。
连续 ACTION 步会在同一次提交后依次执行,上限 32 跳,超出抛 `IllegalStateException`。
连续 ACTION 步会在同一次提交后依次执行,上限 32 跳(**每个 PARALLEL 分支各自计数**),超出抛 `IllegalStateException`。
### 5.1 结构化并行(PARALLEL)
主图仍是单线。`kind: PARALLEL` 是一个步骤,块内多条分支同时推进;全部完成后走该步在父图上的出边。不是 BPMN fork/join 网关。
JSON 嵌套写,引擎拍平存储。分支内边 `to: null` 表示**该分支完成并等待 join**,不是实例通过。PARALLEL 步自己的出边 `to: null` 才结束实例。
v1:至少 2 条分支;禁止套娃 PARALLEL;join 固定 ALL;任一分支按现有规则否决则整单 `REJECTED` 并 SKIPPED 其余 PENDING;`due.goto` 只能指向同一分支内的步骤。会签仍用单步 `candidates` + `ANY`/`ALL`。
事件:`PARALLEL_ENTERED`、`BRANCH_COMPLETED`(`detail` 为 branch id)、`PARALLEL_JOINED`。
## 6. 条件路由
@@ -382,7 +394,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
## 12. 存储
Flyway 脚本按方言分目录:`db/postgresql/migration`、`db/mysql/migration`(各一份当前 schema 的 `V1__baseline.sql`)。未设置 `spring.flyway.locations` 时,starter 按探测到的方言指向对应目录。已有 `flyway_schema_history` 的开发库需清空后重跑。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤/候选人/转移、实例(含 `definition_version`)、任务、`ordo_process_event`、`ordo_action_execution`。
Flyway 脚本按方言分目录:`db/postgresql/migration`、`db/mysql/migration`。未设置 `spring.flyway.locations` 时,starter 按探测到的方言指向对应目录。已有 `flyway_schema_history` 的开发库按版本追加迁移(PARALLEL 为 `V3__parallel_blocks.sql`)。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤(含 `parent_step_id` / `branch_id`)/候选人/转移、实例(含 `definition_version`)、任务、并行令牌 `ordo_instance_token`、`ordo_process_event`、`ordo_action_execution`。
自定义方言:实现 `SqlDialect` 并用 `META-INF/services` 注册,或提供 `SqlDialect` Bean。新增列时每个已支持方言目录各加一条迁移。
@@ -7,10 +7,11 @@ import java.util.Set;
/**
* A named step in a process definition. Approval steps have one or more candidate assignees;
* action steps have an {@code actionKey} invoked by the host {@link ActionHandler}.
* action steps have an {@code actionKey} invoked by the host {@link ActionHandler};
* parallel steps have two or more {@link ParallelBranch}es.
*/
public record ApprovalStep(String id, String name, List<String> candidates, ApprovalPolicy policy, StepKind kind,
String actionKey, StepDue due) {
String actionKey, StepDue due, List<ParallelBranch> branches) {
public ApprovalStep {
Texts.requireText(id, "step id");
Texts.requireText(name, "step name");
@@ -18,6 +19,7 @@ public record ApprovalStep(String id, String name, List<String> candidates, Appr
Objects.requireNonNull(policy, "policy must not be null");
Objects.requireNonNull(candidates, "candidates must not be null");
candidates = List.copyOf(candidates);
branches = branches == null ? List.of() : List.copyOf(branches);
if (kind == StepKind.ACTION) {
if (!candidates.isEmpty()) {
throw new IllegalArgumentException("an action step must not have candidates");
@@ -27,11 +29,47 @@ public record ApprovalStep(String id, String name, List<String> candidates, Appr
if (due != null) {
throw new IllegalArgumentException("an action step must not have due");
}
if (!branches.isEmpty()) {
throw new IllegalArgumentException("an action step must not have branches");
}
} else if (kind == StepKind.PARALLEL) {
if (!candidates.isEmpty()) {
throw new IllegalArgumentException("a parallel step must not have candidates");
}
if (actionKey != null && !actionKey.isBlank()) {
throw new IllegalArgumentException("a parallel step must not have an action key");
}
actionKey = null;
if (due != null) {
throw new IllegalArgumentException("a parallel step must not have due");
}
if (branches.size() < 2) {
throw new IllegalArgumentException("a parallel step must have at least two branches");
}
Set<String> branchIds = new HashSet<>();
Set<String> nestedSteps = new HashSet<>();
for (ParallelBranch branch : branches) {
Objects.requireNonNull(branch, "branch must not be null");
if (!branchIds.add(branch.id())) {
throw new IllegalArgumentException("duplicate branch in step " + id + ": " + branch.id());
}
for (String stepId : branch.stepIds()) {
if (!nestedSteps.add(stepId)) {
throw new IllegalArgumentException("duplicate step across branches in " + id + ": " + stepId);
}
if (stepId.equals(id)) {
throw new IllegalArgumentException("a parallel step must not include itself");
}
}
}
} else {
if (actionKey != null && !actionKey.isBlank()) {
throw new IllegalArgumentException("an approval step must not have an action key");
}
actionKey = null;
if (!branches.isEmpty()) {
throw new IllegalArgumentException("an approval step must not have branches");
}
if (candidates.isEmpty()) {
throw new IllegalArgumentException("a step must have at least one candidate");
}
@@ -46,12 +84,17 @@ public record ApprovalStep(String id, String name, List<String> candidates, Appr
}
public ApprovalStep(String id, String name, List<String> candidates, ApprovalPolicy policy) {
this(id, name, candidates, policy, StepKind.APPROVAL, null, null);
this(id, name, candidates, policy, StepKind.APPROVAL, null, null, List.of());
}
public ApprovalStep(String id, String name, List<String> candidates, ApprovalPolicy policy, StepKind kind,
String actionKey) {
this(id, name, candidates, policy, kind, actionKey, null);
this(id, name, candidates, policy, kind, actionKey, null, List.of());
}
public ApprovalStep(String id, String name, List<String> candidates, ApprovalPolicy policy, StepKind kind,
String actionKey, StepDue due) {
this(id, name, candidates, policy, kind, actionKey, due, List.of());
}
/** Convenience factory for the common case of a single, fixed approver. */
@@ -62,4 +105,8 @@ public record ApprovalStep(String id, String name, List<String> candidates, Appr
public static ApprovalStep action(String id, String name, String actionKey) {
return new ApprovalStep(id, name, List.of(), ApprovalPolicy.ANY, StepKind.ACTION, actionKey);
}
public static ApprovalStep parallel(String id, String name, List<ParallelBranch> branches) {
return new ApprovalStep(id, name, List.of(), ApprovalPolicy.ANY, StepKind.PARALLEL, null, null, branches);
}
}
@@ -0,0 +1,16 @@
package com.jetlumen.ordo.api;
import java.util.Objects;
/** Runtime token for one branch of a {@link StepKind#PARALLEL} block. */
public record InstanceToken(String id, String instanceId, String parallelStepId, String branchId, String currentStepId,
TokenStatus status) {
public InstanceToken {
Texts.requireText(id, "token id");
Texts.requireText(instanceId, "instance id");
Texts.requireText(parallelStepId, "parallel step id");
Texts.requireText(branchId, "branch id");
Texts.requireText(currentStepId, "current step id");
Objects.requireNonNull(status, "status must not be null");
}
}
@@ -0,0 +1,23 @@
package com.jetlumen.ordo.api;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
/** One branch inside a {@link StepKind#PARALLEL} block. The first step id is the branch start. */
public record ParallelBranch(String id, List<String> stepIds) {
public ParallelBranch {
Texts.requireText(id, "branch id");
if (stepIds == null || stepIds.isEmpty()) {
throw new IllegalArgumentException("branch " + id + " must have at least one step");
}
stepIds = List.copyOf(stepIds);
Set<String> distinct = new HashSet<>();
for (String stepId : stepIds) {
Texts.requireText(stepId, "branch step id");
if (!distinct.add(stepId)) {
throw new IllegalArgumentException("duplicate step in branch " + id + ": " + stepId);
}
}
}
}
@@ -0,0 +1,5 @@
package com.jetlumen.ordo.api;
/** Locates an inner step inside a {@link StepKind#PARALLEL} block. */
public record ParallelMembership(ApprovalStep parallel, ParallelBranch branch) {
}
@@ -2,9 +2,12 @@ package com.jetlumen.ordo.api;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
/** Immutable blueprint for an approval process with explicit step transitions. */
@@ -26,6 +29,34 @@ public record ProcessDefinition(String id, int version, String name, List<Approv
throw new IllegalArgumentException("duplicate step id: " + step.id());
}
}
Map<String, String> innerToParallel = new HashMap<>();
Map<String, String> innerToBranch = new HashMap<>();
Set<String> innerIds = new HashSet<>();
for (ApprovalStep step : steps) {
if (step.kind() != StepKind.PARALLEL) {
continue;
}
for (ParallelBranch branch : step.branches()) {
for (String nestedId : branch.stepIds()) {
if (!ids.contains(nestedId)) {
throw new IllegalArgumentException("unknown branch step id: " + nestedId);
}
ApprovalStep nested = stepById(steps, nestedId);
if (nested.kind() == StepKind.PARALLEL) {
throw new IllegalArgumentException("nested parallel is not allowed: " + nestedId);
}
String previous = innerToParallel.put(nestedId, step.id());
if (previous != null) {
throw new IllegalArgumentException("step belongs to multiple parallel blocks: " + nestedId);
}
innerToBranch.put(nestedId, branch.id());
innerIds.add(nestedId);
}
}
}
if (innerIds.contains(steps.get(0).id())) {
throw new IllegalArgumentException("start step must not be a parallel branch step");
}
Objects.requireNonNull(transitions, "transitions must not be null");
Set<String> stepsWithOutgoing = new HashSet<>();
Set<String> fromPriorityKeys = new HashSet<>();
@@ -36,6 +67,19 @@ public record ProcessDefinition(String id, int version, String name, List<Approv
if (transition.toStepId() != null && !ids.contains(transition.toStepId())) {
throw new IllegalArgumentException("unknown toStepId: " + transition.toStepId());
}
String fromParallel = innerToParallel.get(transition.fromStepId());
if (fromParallel != null) {
if (transition.toStepId() != null) {
String toParallel = innerToParallel.get(transition.toStepId());
if (toParallel == null || !fromParallel.equals(toParallel)
|| !innerToBranch.get(transition.fromStepId()).equals(innerToBranch.get(transition.toStepId()))) {
throw new IllegalArgumentException(
"branch transition must stay in the same branch: " + transition.fromStepId());
}
}
} else if (transition.toStepId() != null && innerIds.contains(transition.toStepId())) {
throw new IllegalArgumentException("parent graph must not target a branch step: " + transition.toStepId());
}
stepsWithOutgoing.add(transition.fromStepId());
String fromPriority = transition.fromStepId() + '\0' + transition.priority();
if (!fromPriorityKeys.add(fromPriority)) {
@@ -47,8 +91,21 @@ public record ProcessDefinition(String id, int version, String name, List<Approv
if (!stepsWithOutgoing.contains(step.id())) {
throw new IllegalArgumentException("step " + step.id() + " has no outgoing transition");
}
if (step.due() != null && step.due().then() == DueThen.GOTO && !ids.contains(step.due().to())) {
throw new IllegalArgumentException("unknown due toStepId: " + step.due().to());
if (step.due() != null && step.due().then() == DueThen.GOTO) {
String target = step.due().to();
if (!ids.contains(target)) {
throw new IllegalArgumentException("unknown due toStepId: " + target);
}
String fromParallel = innerToParallel.get(step.id());
if (fromParallel != null) {
String toParallel = innerToParallel.get(target);
if (toParallel == null || !fromParallel.equals(toParallel)
|| !innerToBranch.get(step.id()).equals(innerToBranch.get(target))) {
throw new IllegalArgumentException("due goto must stay in the same branch: " + step.id());
}
} else if (innerIds.contains(target)) {
throw new IllegalArgumentException("due goto must not target a branch step: " + target);
}
}
}
transitions = transitions.stream()
@@ -74,11 +131,30 @@ public record ProcessDefinition(String id, int version, String name, List<Approv
&& transitions.equals(other.transitions);
}
public Optional<ParallelMembership> membershipOf(String stepId) {
for (ApprovalStep step : steps) {
if (step.kind() != StepKind.PARALLEL) {
continue;
}
for (ParallelBranch branch : step.branches()) {
if (branch.stepIds().contains(stepId)) {
return Optional.of(new ParallelMembership(step, branch));
}
}
}
return Optional.empty();
}
/**
* Builds a definition whose transitions mirror the former linear steps order: each step
* unconditionally advances to the next, and the last step unconditionally ends.
*/
public static ProcessDefinition linear(String id, String name, List<ApprovalStep> steps) {
for (ApprovalStep step : steps) {
if (step.kind() == StepKind.PARALLEL) {
throw new IllegalArgumentException("linear definitions cannot contain parallel steps");
}
}
List<StepTransition> transitions = new ArrayList<>();
for (int i = 0; i < steps.size() - 1; i++) {
transitions.add(StepTransition.always(steps.get(i).id(), steps.get(i + 1).id()));
@@ -88,4 +164,13 @@ public record ProcessDefinition(String id, int version, String name, List<Approv
}
return new ProcessDefinition(id, name, steps, transitions);
}
private static ApprovalStep stepById(List<ApprovalStep> steps, String stepId) {
for (ApprovalStep step : steps) {
if (step.id().equals(stepId)) {
return step;
}
}
throw new IllegalArgumentException("unknown branch step id: " + stepId);
}
}
@@ -10,9 +10,11 @@ import java.io.InputStream;
import java.time.DateTimeException;
import java.time.Duration;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Locale;
import java.util.Objects;
import java.util.Set;
/** Parses a structural JSON process graph into a {@link ProcessDefinition}. */
public final class ProcessDefinitionParser {
@@ -52,51 +54,149 @@ public final class ProcessDefinitionParser {
if (document.steps() == null || document.steps().isEmpty()) {
throw new IllegalArgumentException("a definition must contain at least one step");
}
List<ApprovalStep> steps = new ArrayList<>(document.steps().size());
List<ApprovalStep> steps = new ArrayList<>();
List<StepTransition> transitions = new ArrayList<>();
for (StepDocument step : document.steps()) {
if (step == null) {
throw new IllegalArgumentException("step must not be null");
}
ApprovalPolicy policy = step.policy() == null ? ApprovalPolicy.ANY : step.policy();
StepKind kind = step.kind() == null ? StepKind.APPROVAL : step.kind();
List<String> candidates = step.candidates() == null ? List.of() : step.candidates();
steps.add(new ApprovalStep(step.id(), step.name(), candidates, policy, kind, step.action(),
parseDue(step.due())));
flattenStep(step, steps, transitions, false);
}
rotateStartStep(steps, document.startStep());
List<StepTransition> transitions = new ArrayList<>();
if (document.transitions() != null) {
for (TransitionDocument transition : document.transitions()) {
if (transition == null) {
throw new IllegalArgumentException("transition must not be null");
}
int priority = transition.priority() == null ? 0 : transition.priority();
transitions.add(new StepTransition(transition.from(), transition.to(),
RoutingWhen.parse(transition.when()), priority));
transitions.add(toTransition(transition));
}
}
return new ProcessDefinition(document.id(), document.name(), steps, transitions);
}
private static void flattenStep(StepDocument step, List<ApprovalStep> steps, List<StepTransition> transitions,
boolean nested) {
StepKind kind = step.kind() == null ? StepKind.APPROVAL : step.kind();
if (kind == StepKind.PARALLEL) {
if (nested) {
throw new IllegalArgumentException("nested parallel is not allowed: " + step.id());
}
if (step.branches() == null || step.branches().size() < 2) {
throw new IllegalArgumentException("a parallel step must have at least two branches");
}
List<ParallelBranch> branches = new ArrayList<>();
List<ApprovalStep> nestedSteps = new ArrayList<>();
for (BranchDocument branch : step.branches()) {
if (branch == null) {
throw new IllegalArgumentException("branch must not be null");
}
if (branch.steps() == null || branch.steps().isEmpty()) {
throw new IllegalArgumentException("branch " + branch.id() + " must have at least one step");
}
List<String> stepIds = new ArrayList<>();
for (StepDocument nestedStep : branch.steps()) {
if (nestedStep == null) {
throw new IllegalArgumentException("step must not be null");
}
flattenStep(nestedStep, nestedSteps, transitions, true);
stepIds.add(nestedStep.id());
}
if (branch.transitions() != null) {
for (TransitionDocument transition : branch.transitions()) {
if (transition == null) {
throw new IllegalArgumentException("transition must not be null");
}
transitions.add(toTransition(transition));
}
}
branches.add(new ParallelBranch(branch.id(), stepIds));
}
steps.add(ApprovalStep.parallel(step.id(), step.name(), branches));
steps.addAll(nestedSteps);
return;
}
ApprovalPolicy policy = step.policy() == null ? ApprovalPolicy.ANY : step.policy();
List<String> candidates = step.candidates() == null ? List.of() : step.candidates();
steps.add(new ApprovalStep(step.id(), step.name(), candidates, policy, kind, step.action(),
parseDue(step.due())));
}
private static StepTransition toTransition(TransitionDocument transition) {
int priority = transition.priority() == null ? 0 : transition.priority();
return new StepTransition(transition.from(), transition.to(), RoutingWhen.parse(transition.when()), priority);
}
private static DefinitionDocument fromDefinition(ProcessDefinition definition) {
List<StepDocument> steps = new ArrayList<>(definition.steps().size());
Set<String> innerIds = new HashSet<>();
for (ApprovalStep step : definition.steps()) {
if (step.kind() != StepKind.PARALLEL) {
continue;
}
for (ParallelBranch branch : step.branches()) {
innerIds.addAll(branch.stepIds());
}
}
List<StepDocument> steps = new ArrayList<>();
for (ApprovalStep step : definition.steps()) {
if (innerIds.contains(step.id())) {
continue;
}
steps.add(toStepDocument(definition, step));
}
List<TransitionDocument> transitions = new ArrayList<>();
for (StepTransition transition : definition.transitions()) {
if (innerIds.contains(transition.fromStepId())) {
continue;
}
transitions.add(toTransitionDocument(transition));
}
return new DefinitionDocument(definition.id(), definition.name(), definition.version(),
definition.steps().get(0).id(), steps, transitions);
}
private static StepDocument toStepDocument(ProcessDefinition definition, ApprovalStep step) {
if (step.kind() == StepKind.PARALLEL) {
List<BranchDocument> branches = new ArrayList<>();
for (ParallelBranch branch : step.branches()) {
List<StepDocument> branchSteps = new ArrayList<>();
for (String stepId : branch.stepIds()) {
branchSteps.add(toLeafDocument(requireStep(definition, stepId)));
}
List<TransitionDocument> branchTransitions = new ArrayList<>();
for (StepTransition transition : definition.transitions()) {
if (branch.stepIds().contains(transition.fromStepId())) {
branchTransitions.add(toTransitionDocument(transition));
}
}
branches.add(new BranchDocument(branch.id(), branchSteps, branchTransitions));
}
return new StepDocument(step.id(), step.name(), null, null, StepKind.PARALLEL, null, null, branches);
}
return toLeafDocument(step);
}
private static StepDocument toLeafDocument(ApprovalStep step) {
DueDocument due = null;
if (step.due() != null) {
StepDue stepDue = step.due();
due = new DueDocument(stepDue.after().toString(), stepDue.then().name().toLowerCase(Locale.ROOT),
stepDue.to(), stepDue.action());
}
steps.add(new StepDocument(step.id(), step.name(), step.candidates(), step.policy(), step.kind(),
step.actionKey(), due));
List<String> candidates = step.candidates().isEmpty() ? null : step.candidates();
return new StepDocument(step.id(), step.name(), candidates, step.policy(), step.kind(),
step.actionKey(), due, null);
}
List<TransitionDocument> transitions = new ArrayList<>(definition.transitions().size());
for (StepTransition transition : definition.transitions()) {
transitions.add(new TransitionDocument(transition.fromStepId(), transition.toStepId(),
transition.when() == null ? null : transition.when().toTree(), transition.priority()));
private static TransitionDocument toTransitionDocument(StepTransition transition) {
return new TransitionDocument(transition.fromStepId(), transition.toStepId(),
transition.when() == null ? null : transition.when().toTree(), transition.priority());
}
return new DefinitionDocument(definition.id(), definition.name(), definition.version(),
definition.steps().get(0).id(), steps, transitions);
private static ApprovalStep requireStep(ProcessDefinition definition, String stepId) {
return definition.steps().stream()
.filter(step -> step.id().equals(stepId))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("unknown step: " + stepId));
}
private static void rotateStartStep(List<ApprovalStep> steps, String startStepId) {
@@ -158,7 +258,14 @@ public final class ProcessDefinitionParser {
ApprovalPolicy policy,
StepKind kind,
@JsonProperty("action") String action,
DueDocument due) {
DueDocument due,
List<BranchDocument> branches) {
}
private record BranchDocument(
String id,
List<StepDocument> steps,
List<TransitionDocument> transitions) {
}
private record DueDocument(
@@ -14,5 +14,8 @@ public enum ProcessEventType {
INSTANCE_WITHDRAWN,
INSTANCE_CANCELLED,
ACTION_SUCCEEDED,
ACTION_FAILED
ACTION_FAILED,
PARALLEL_ENTERED,
BRANCH_COMPLETED,
PARALLEL_JOINED
}
@@ -1,5 +1,5 @@
package com.jetlumen.ordo.api;
public enum StepKind {
APPROVAL, ACTION
APPROVAL, ACTION, PARALLEL
}
@@ -0,0 +1,6 @@
package com.jetlumen.ordo.api;
/** Status of a parallel-branch token. */
public enum TokenStatus {
ACTIVE, COMPLETED
}
@@ -0,0 +1,20 @@
package com.jetlumen.ordo.api.repository;
import com.jetlumen.ordo.api.InstanceToken;
import java.util.List;
/** Storage port for parallel-branch tokens. */
public interface InstanceTokenRepository {
void insert(InstanceToken token);
List<InstanceToken> findByInstanceId(String instanceId);
List<InstanceToken> findByInstanceAndParallel(String instanceId, String parallelStepId);
boolean completeIfActive(String instanceId, String parallelStepId, String branchId);
boolean updateCurrentStepIfActive(String instanceId, String parallelStepId, String branchId, String currentStepId);
void completeActiveByInstanceId(String instanceId);
}
@@ -267,4 +267,146 @@ class ProcessDefinitionParserTest {
""";
assertThrows(IllegalArgumentException.class, () -> ProcessDefinitionParser.fromJson(badThen));
}
@Test
void roundTripsAParallelBlock() {
String json = """
{
"id": "p",
"name": "P",
"steps": [
{
"id": "dept",
"name": "Dept",
"kind": "PARALLEL",
"branches": [
{
"id": "legal",
"steps": [
{ "id": "legal", "name": "Legal", "candidates": ["lee"] }
],
"transitions": [{ "from": "legal", "to": null }]
},
{
"id": "finance",
"steps": [
{ "id": "finance", "name": "Finance", "candidates": ["fay"] }
],
"transitions": [{ "from": "finance", "to": null }]
}
]
},
{ "id": "ceo", "name": "CEO", "candidates": ["cara"] }
],
"transitions": [
{ "from": "dept", "to": "ceo" },
{ "from": "ceo", "to": null }
]
}
""";
ProcessDefinition definition = ProcessDefinitionParser.fromJson(json);
assertEquals(StepKind.PARALLEL, definition.steps().get(0).kind());
ProcessDefinition roundTrip = ProcessDefinitionParser.fromJson(ProcessDefinitionParser.toJson(definition));
assertTrue(definition.sameGraph(roundTrip));
}
@Test
void rejectsParallelWithFewerThanTwoBranchesOrNestingOrCrossBranch() {
String oneBranch = """
{
"id": "p", "name": "P",
"steps": [{
"id": "dept", "name": "Dept", "kind": "PARALLEL",
"branches": [{
"id": "legal",
"steps": [{ "id": "legal", "name": "Legal", "candidates": ["lee"] }],
"transitions": [{ "from": "legal", "to": null }]
}]
}],
"transitions": [{ "from": "dept", "to": null }]
}
""";
assertThrows(IllegalArgumentException.class, () -> ProcessDefinitionParser.fromJson(oneBranch));
String nested = """
{
"id": "p", "name": "P",
"steps": [{
"id": "dept", "name": "Dept", "kind": "PARALLEL",
"branches": [
{
"id": "legal",
"steps": [{
"id": "inner", "name": "Inner", "kind": "PARALLEL",
"branches": [
{ "id": "a", "steps": [{ "id": "a", "name": "A", "candidates": ["a"] }],
"transitions": [{ "from": "a", "to": null }] },
{ "id": "b", "steps": [{ "id": "b", "name": "B", "candidates": ["b"] }],
"transitions": [{ "from": "b", "to": null }] }
]
}],
"transitions": [{ "from": "inner", "to": null }]
},
{
"id": "finance",
"steps": [{ "id": "finance", "name": "Finance", "candidates": ["fay"] }],
"transitions": [{ "from": "finance", "to": null }]
}
]
}],
"transitions": [{ "from": "dept", "to": null }]
}
""";
assertThrows(IllegalArgumentException.class, () -> ProcessDefinitionParser.fromJson(nested));
String cross = """
{
"id": "p", "name": "P",
"steps": [{
"id": "dept", "name": "Dept", "kind": "PARALLEL",
"branches": [
{
"id": "legal",
"steps": [{ "id": "legal", "name": "Legal", "candidates": ["lee"] }],
"transitions": [{ "from": "legal", "to": "finance" }]
},
{
"id": "finance",
"steps": [{ "id": "finance", "name": "Finance", "candidates": ["fay"] }],
"transitions": [{ "from": "finance", "to": null }]
}
]
}],
"transitions": [{ "from": "dept", "to": null }]
}
""";
assertThrows(IllegalArgumentException.class, () -> ProcessDefinitionParser.fromJson(cross));
}
@Test
void rejectsDueGotoOutOfBranchInJson() {
String json = """
{
"id": "p", "name": "P",
"steps": [{
"id": "dept", "name": "Dept", "kind": "PARALLEL",
"branches": [
{
"id": "legal",
"steps": [{
"id": "legal", "name": "Legal", "candidates": ["lee"],
"due": { "after": "PT1H", "then": "goto", "to": "finance" }
}],
"transitions": [{ "from": "legal", "to": null }]
},
{
"id": "finance",
"steps": [{ "id": "finance", "name": "Finance", "candidates": ["fay"] }],
"transitions": [{ "from": "finance", "to": null }]
}
]
}],
"transitions": [{ "from": "dept", "to": null }]
}
""";
assertThrows(IllegalArgumentException.class, () -> ProcessDefinitionParser.fromJson(json));
}
}
@@ -133,4 +133,76 @@ class ProcessDefinitionTest {
null, StepDue.gotoStep(java.time.Duration.ofHours(1), "missing"))),
List.of(StepTransition.end("manager"))));
}
@Test
void rejectsLinearDefinitionsThatContainParallelSteps() {
assertThrows(IllegalArgumentException.class, () -> ProcessDefinition.linear("leave", "Leave", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("a", List.of("a")),
new ParallelBranch("b", List.of("b")))))));
}
@Test
void rejectsNestedParallelAndCrossBranchEdges() {
ApprovalStep dept = ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance"))));
assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition("p", "P", List.of(
dept,
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay")),
List.of(
StepTransition.always("legal", "finance"),
StepTransition.end("finance"),
StepTransition.end("dept"))));
assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition("p", "P", List.of(
dept,
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.parallel("finance", "Finance", List.of(
new ParallelBranch("x", List.of("x1")),
new ParallelBranch("y", List.of("y1")))),
ApprovalStep.single("x1", "X", "x"),
ApprovalStep.single("y1", "Y", "y")),
List.of(
StepTransition.end("legal"),
StepTransition.end("x1"),
StepTransition.end("y1"),
StepTransition.end("finance"),
StepTransition.end("dept"))));
}
@Test
void rejectsDueGotoLeavingABranch() {
assertThrows(IllegalArgumentException.class, () -> new ProcessDefinition("p", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
new ApprovalStep("legal", "Legal", List.of("lee"), ApprovalPolicy.ANY, StepKind.APPROVAL, null,
StepDue.gotoStep(java.time.Duration.ofHours(1), "finance")),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo"))));
}
@Test
void acceptsAStructuredParallelBlock() {
ProcessDefinition definition = new ProcessDefinition("p", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo")));
assertEquals(StepKind.PARALLEL, definition.steps().get(0).kind());
assertEquals("legal", definition.membershipOf("legal").orElseThrow().branch().id());
}
}
@@ -35,8 +35,13 @@ import com.jetlumen.ordo.api.query.InstanceQuery;
import com.jetlumen.ordo.api.query.Page;
import com.jetlumen.ordo.api.query.PageRequest;
import com.jetlumen.ordo.api.query.TaskQuery;
import com.jetlumen.ordo.api.InstanceToken;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ParallelMembership;
import com.jetlumen.ordo.api.TokenStatus;
import com.jetlumen.ordo.api.repository.ActionExecutionRepository;
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
import com.jetlumen.ordo.api.repository.InstanceTokenRepository;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import com.jetlumen.ordo.api.repository.ProcessHistoryRepository;
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
@@ -75,6 +80,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
private final ApprovalTaskRepository taskRepository;
private final ProcessHistoryRepository historyRepository;
private final ActionExecutionRepository actionExecutionRepository;
private final InstanceTokenRepository tokenRepository;
private final List<OrdoEventListener> listeners;
/** Monotonic millis offset so same-clock events keep causal order after JDBC Timestamp round-trips. */
private final AtomicLong eventSequence = new AtomicLong();
@@ -86,6 +92,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
ApprovalTaskRepository taskRepository,
ProcessHistoryRepository historyRepository,
ActionExecutionRepository actionExecutionRepository,
InstanceTokenRepository tokenRepository,
List<OrdoEventListener> listeners) {
this.clock = Objects.requireNonNull(clock, "clock must not be null");
this.assigneeResolver = Objects.requireNonNull(assigneeResolver, "assigneeResolver must not be null");
@@ -98,6 +105,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
this.historyRepository = Objects.requireNonNull(historyRepository, "historyRepository must not be null");
this.actionExecutionRepository = Objects.requireNonNull(actionExecutionRepository,
"actionExecutionRepository must not be null");
this.tokenRepository = Objects.requireNonNull(tokenRepository, "tokenRepository must not be null");
this.listeners = List.copyOf(Objects.requireNonNull(listeners, "listeners must not be null"));
}
@@ -222,16 +230,8 @@ public final class DefaultOrdoEngine implements OrdoEngine {
throw new InstanceAlreadyCompletedException(instanceId);
}
record(events, instanceId, null, null, eventType, actor, comment, now);
TaskAction action = new TaskAction(actor, comment, now);
for (ApprovalTask pending : taskRepository.findPendingByInstanceId(instanceId)) {
ApprovalTask skipped = new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(),
pending.name(), pending.assignee(), TaskStatus.SKIPPED, pending.createdAt(), now, action,
pending.dueAt());
if (taskRepository.completeIfPending(skipped)) {
record(events, instanceId, pending.id(), pending.stepId(), ProcessEventType.TASK_SKIPPED, actor,
comment, now);
}
}
skipRemainingPending(instanceId, actor, comment, now, events);
tokenRepository.completeActiveByInstanceId(instanceId);
return finished;
});
dispatch(events);
@@ -409,6 +409,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
List<ApprovalTask> siblings = taskRepository.findByInstanceIdAndStepId(instance.id(), overdue.stepId());
skipPendingSiblings(siblings, "", null, now, events);
enterStep(instance, definition, requireStep(definition, due.to()), now, queued, events);
touchToken(definition, instance.id(), due.to());
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
due.to(), now);
return true;
@@ -498,6 +499,9 @@ public final class DefaultOrdoEngine implements OrdoEngine {
Instant now, List<PendingAction> queued, List<ProcessEvent> events) {
StepTransition matched = resolveTransition(definition, step, instance);
if (matched.toStepId() == null) {
if (completeBranchIfInner(instance, definition, step, now, queued, events)) {
return;
}
completeInstance(instance, ProcessStatus.APPROVED, now, events);
} else {
enterStep(instance, definition, requireStep(definition, matched.toStepId()), now, queued, events);
@@ -508,8 +512,13 @@ public final class DefaultOrdoEngine implements OrdoEngine {
List<PendingAction> queued, List<ProcessEvent> events) {
ApprovalStep current = start;
for (int hops = 0; hops < MAX_CONSECUTIVE_ACTIONS; hops++) {
if (current.kind() == StepKind.PARALLEL) {
enterParallel(instance, definition, current, now, queued, events);
return;
}
if (current.kind() == StepKind.APPROVAL) {
createStepTasks(instance, current, now, events);
touchToken(definition, instance.id(), current.id());
return;
}
String executionId = nextId();
@@ -518,8 +527,12 @@ public final class DefaultOrdoEngine implements OrdoEngine {
now.plusMillis(eventSequence.getAndIncrement()), null));
queued.add(new PendingAction(executionId, current.actionKey(), instance.id(), current.id(),
instance.context()));
touchToken(definition, instance.id(), current.id());
StepTransition matched = resolveTransition(definition, current, instance);
if (matched.toStepId() == null) {
if (completeBranchIfInner(instance, definition, current, now, queued, events)) {
return;
}
completeInstance(instance, ProcessStatus.APPROVED, now, events);
return;
}
@@ -528,6 +541,44 @@ public final class DefaultOrdoEngine implements OrdoEngine {
throw new IllegalStateException("too many consecutive action steps in instance: " + instance.id());
}
private void enterParallel(ProcessInstance instance, ProcessDefinition definition, ApprovalStep parallel,
Instant now, List<PendingAction> queued, List<ProcessEvent> events) {
record(events, instance.id(), null, parallel.id(), ProcessEventType.PARALLEL_ENTERED, null, parallel.id(), now);
for (var branch : parallel.branches()) {
String startId = branch.stepIds().get(0);
tokenRepository.insert(new InstanceToken(nextId(), instance.id(), parallel.id(), branch.id(), startId,
TokenStatus.ACTIVE));
enterStep(instance, definition, requireStep(definition, startId), now, queued, events);
}
}
private boolean completeBranchIfInner(ProcessInstance instance, ProcessDefinition definition, ApprovalStep step,
Instant now, List<PendingAction> queued, List<ProcessEvent> events) {
ParallelMembership membership = definition.membershipOf(step.id()).orElse(null);
if (membership == null) {
return false;
}
ApprovalStep parallel = membership.parallel();
ParallelBranch branch = membership.branch();
tokenRepository.completeIfActive(instance.id(), parallel.id(), branch.id());
record(events, instance.id(), null, parallel.id(), ProcessEventType.BRANCH_COMPLETED, null, branch.id(), now);
List<InstanceToken> tokens = tokenRepository.findByInstanceAndParallel(instance.id(), parallel.id());
boolean joined = tokens.size() == parallel.branches().size()
&& tokens.stream().allMatch(token -> token.status() == TokenStatus.COMPLETED);
if (joined) {
record(events, instance.id(), null, parallel.id(), ProcessEventType.PARALLEL_JOINED, null, parallel.id(),
now);
advanceOrComplete(instance, definition, parallel, now, queued, events);
}
return true;
}
private void touchToken(ProcessDefinition definition, String instanceId, String stepId) {
definition.membershipOf(stepId).ifPresent(membership ->
tokenRepository.updateCurrentStepIfActive(instanceId, membership.parallel().id(),
membership.branch().id(), stepId));
}
private void finishCommittedWork(List<PendingAction> queued, List<ProcessEvent> events) {
runQueuedActions(queued, events);
dispatch(events);
@@ -625,6 +676,22 @@ public final class DefaultOrdoEngine implements OrdoEngine {
ProcessEventType type = status == ProcessStatus.APPROVED
? ProcessEventType.INSTANCE_APPROVED : ProcessEventType.INSTANCE_REJECTED;
record(events, instance.id(), null, null, type, null, null, now);
skipRemainingPending(instance.id(), null, null, now, events);
tokenRepository.completeActiveByInstanceId(instance.id());
}
private void skipRemainingPending(String instanceId, String actor, String comment, Instant now,
List<ProcessEvent> events) {
TaskAction action = actor == null ? null : new TaskAction(actor, comment, now);
for (ApprovalTask pending : taskRepository.findPendingByInstanceId(instanceId)) {
ApprovalTask skipped = new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(),
pending.name(), pending.assignee(), TaskStatus.SKIPPED, pending.createdAt(), now, action,
pending.dueAt());
if (taskRepository.completeIfPending(skipped)) {
record(events, instanceId, pending.id(), pending.stepId(), ProcessEventType.TASK_SKIPPED, actor,
comment, now);
}
}
}
private ProcessDefinition requireLatestDefinition(String definitionId) {
@@ -17,6 +17,7 @@ import com.jetlumen.ordo.api.query.PageRequest;
import com.jetlumen.ordo.api.query.TaskQuery;
import com.jetlumen.ordo.core.repository.InMemoryActionExecutionRepository;
import com.jetlumen.ordo.core.repository.InMemoryApprovalTaskRepository;
import com.jetlumen.ordo.core.repository.InMemoryInstanceTokenRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessHistoryRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository;
@@ -72,6 +73,7 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
new InMemoryApprovalTaskRepository(instanceRepository),
new InMemoryProcessHistoryRepository(),
new InMemoryActionExecutionRepository(),
new InMemoryInstanceTokenRepository(),
listeners);
}
@@ -0,0 +1,77 @@
package com.jetlumen.ordo.core.repository;
import com.jetlumen.ordo.api.InstanceToken;
import com.jetlumen.ordo.api.TokenStatus;
import com.jetlumen.ordo.api.repository.InstanceTokenRepository;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
/** Development-only in-memory implementation of the parallel token port. */
public final class InMemoryInstanceTokenRepository implements InstanceTokenRepository {
private final Map<String, InstanceToken> tokens = new LinkedHashMap<>();
@Override
public synchronized void insert(InstanceToken token) {
Objects.requireNonNull(token, "token must not be null");
if (tokens.putIfAbsent(key(token.instanceId(), token.parallelStepId(), token.branchId()), token) != null) {
throw new IllegalStateException("token already exists for branch: " + token.branchId());
}
}
@Override
public synchronized List<InstanceToken> findByInstanceId(String instanceId) {
return tokens.values().stream().filter(token -> token.instanceId().equals(instanceId)).toList();
}
@Override
public synchronized List<InstanceToken> findByInstanceAndParallel(String instanceId, String parallelStepId) {
return tokens.values().stream()
.filter(token -> token.instanceId().equals(instanceId) && token.parallelStepId().equals(parallelStepId))
.toList();
}
@Override
public synchronized boolean completeIfActive(String instanceId, String parallelStepId, String branchId) {
String key = key(instanceId, parallelStepId, branchId);
InstanceToken current = tokens.get(key);
if (current == null || current.status() != TokenStatus.ACTIVE) {
return false;
}
tokens.put(key, new InstanceToken(current.id(), current.instanceId(), current.parallelStepId(),
current.branchId(), current.currentStepId(), TokenStatus.COMPLETED));
return true;
}
@Override
public synchronized boolean updateCurrentStepIfActive(String instanceId, String parallelStepId, String branchId,
String currentStepId) {
String key = key(instanceId, parallelStepId, branchId);
InstanceToken current = tokens.get(key);
if (current == null || current.status() != TokenStatus.ACTIVE) {
return false;
}
tokens.put(key, new InstanceToken(current.id(), current.instanceId(), current.parallelStepId(),
current.branchId(), currentStepId, TokenStatus.ACTIVE));
return true;
}
@Override
public synchronized void completeActiveByInstanceId(String instanceId) {
List<InstanceToken> snapshot = new ArrayList<>(tokens.values());
for (InstanceToken token : snapshot) {
if (token.instanceId().equals(instanceId) && token.status() == TokenStatus.ACTIVE) {
tokens.put(key(token.instanceId(), token.parallelStepId(), token.branchId()),
new InstanceToken(token.id(), token.instanceId(), token.parallelStepId(), token.branchId(),
token.currentStepId(), TokenStatus.COMPLETED));
}
}
}
private static String key(String instanceId, String parallelStepId, String branchId) {
return instanceId + '\0' + parallelStepId + '\0' + branchId;
}
}
@@ -6,6 +6,7 @@ import com.jetlumen.ordo.api.AssigneeResolver;
import com.jetlumen.ordo.api.ApprovalPolicy;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessEvent;
@@ -853,6 +854,128 @@ class InMemoryOrdoEngineTest {
assertThrows(IllegalArgumentException.class, () -> engine.processDue(0));
}
@Test
void parallelJoinsAfterEveryBranchCompletes() {
engine.publish(parallelReview());
var instance = engine.start("parallel-review", "alice");
List<ApprovalTask> pending = engine.findPendingTasksByInstanceId(instance.id());
assertEquals(2, pending.size());
ApprovalTask legal = pending.stream().filter(task -> task.stepId().equals("legal")).findFirst().orElseThrow();
ApprovalTask finance = pending.stream().filter(task -> task.stepId().equals("finance")).findFirst().orElseThrow();
engine.approve(legal.id(), "lee");
assertEquals(ProcessStatus.RUNNING, engine.findInstance(instance.id()).orElseThrow().status());
assertEquals(1, engine.findPendingTasksByInstanceId(instance.id()).size());
assertEquals("finance", engine.findPendingTasksByInstanceId(instance.id()).get(0).stepId());
engine.approve(finance.id(), "fay");
ApprovalTask ceo = engine.findPendingTasksByInstanceId(instance.id()).get(0);
assertEquals("ceo", ceo.stepId());
engine.approve(ceo.id(), "cara");
assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status());
List<ProcessEventType> types = engine.queryHistory(instance.id(), new PageRequest(0, 50)).content().stream()
.map(ProcessEvent::type).toList();
assertTrue(types.contains(ProcessEventType.PARALLEL_ENTERED));
assertTrue(types.contains(ProcessEventType.BRANCH_COMPLETED));
assertTrue(types.contains(ProcessEventType.PARALLEL_JOINED));
}
@Test
void parallelRejectSkipsTheOtherBranch() {
engine.publish(parallelReview());
var instance = engine.start("parallel-review", "alice");
ApprovalTask legal = engine.findPendingTasksByInstanceId(instance.id()).stream()
.filter(task -> task.stepId().equals("legal")).findFirst().orElseThrow();
ApprovalTask finance = engine.findPendingTasksByInstanceId(instance.id()).stream()
.filter(task -> task.stepId().equals("finance")).findFirst().orElseThrow();
engine.reject(legal.id(), "lee");
assertEquals(ProcessStatus.REJECTED, engine.findInstance(instance.id()).orElseThrow().status());
assertEquals(TaskStatus.SKIPPED, engine.findTask(finance.id()).orElseThrow().status());
assertTrue(engine.findPendingTasksByInstanceId(instance.id()).isEmpty());
}
@Test
void parallelWithdrawAndCancelSkipBranchTasks() {
engine.publish(parallelReview());
var withdrawn = engine.start("parallel-review", "alice");
engine.withdraw(withdrawn.id(), "alice");
assertEquals(ProcessStatus.WITHDRAWN, engine.findInstance(withdrawn.id()).orElseThrow().status());
assertTrue(engine.findPendingTasksByInstanceId(withdrawn.id()).isEmpty());
var cancelled = engine.start("parallel-review", "alice");
engine.cancel(cancelled.id(), "admin");
assertEquals(ProcessStatus.CANCELLED, engine.findInstance(cancelled.id()).orElseThrow().status());
assertTrue(engine.findPendingTasksByInstanceId(cancelled.id()).isEmpty());
}
@Test
void parallelBranchCanRunActionThenJoin() {
engine.publish(new ProcessDefinition("parallel-action", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal-mail")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.action("legal-mail", "Legal mail", "ok-mail"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal-mail"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo"))));
var instance = engine.start("parallel-action", "alice");
assertEquals(1, engine.findPendingTasksByInstanceId(instance.id()).size());
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "fay");
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "cara");
assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void processDueGotoStaysInsideABranch() {
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock);
dueEngine.publish(new ProcessDefinition("parallel-due", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal", "legal-2")),
new ParallelBranch("finance", List.of("finance")))),
new ApprovalStep("legal", "Legal", List.of("lee"), ApprovalPolicy.ANY, StepKind.APPROVAL, null,
StepDue.gotoStep(java.time.Duration.ofHours(1), "legal-2")),
ApprovalStep.single("legal-2", "Legal 2", "lisa"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("legal-2"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo"))));
var instance = dueEngine.start("parallel-due", "alice");
clock.set(Instant.parse("2026-01-15T10:00:00Z"));
assertEquals(1, dueEngine.processDue(10));
assertTrue(dueEngine.findPendingTasksByInstanceId(instance.id()).stream()
.anyMatch(task -> task.stepId().equals("legal-2")));
assertTrue(dueEngine.findPendingTasksByInstanceId(instance.id()).stream()
.anyMatch(task -> task.stepId().equals("finance")));
dueEngine.approve(dueEngine.findPendingTasksByAssignee("lisa").get(0).id(), "lisa");
dueEngine.approve(dueEngine.findPendingTasksByAssignee("fay").get(0).id(), "fay");
dueEngine.approve(dueEngine.findPendingTasksByAssignee("cara").get(0).id(), "cara");
assertEquals(ProcessStatus.APPROVED, dueEngine.findInstance(instance.id()).orElseThrow().status());
}
private static ProcessDefinition parallelReview() {
return new ProcessDefinition("parallel-review", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo")));
}
private static final class MutableClock extends Clock {
private Instant instant;
@@ -8,6 +8,7 @@ import com.jetlumen.ordo.api.RoutingCondition;
import com.jetlumen.ordo.api.TransactionExecutor;
import com.jetlumen.ordo.api.repository.ActionExecutionRepository;
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
import com.jetlumen.ordo.api.repository.InstanceTokenRepository;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import com.jetlumen.ordo.api.repository.ProcessHistoryRepository;
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
@@ -16,6 +17,7 @@ import com.jetlumen.ordo.spring.OrdoProperties;
import com.jetlumen.ordo.storage.jdbc.JdbcActionExecutionRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcApprovalTaskRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcConnectionProvider;
import com.jetlumen.ordo.storage.jdbc.JdbcInstanceTokenRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcProcessDefinitionRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcProcessHistoryRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcProcessInstanceRepository;
@@ -134,6 +136,13 @@ public class OrdoJdbcAutoConfiguration {
return new JdbcActionExecutionRepository(connectionProvider, ordoSqlDialect);
}
@Bean
@ConditionalOnMissingBean
public InstanceTokenRepository ordoInstanceTokenRepository(
JdbcConnectionProvider connectionProvider, SqlDialect ordoSqlDialect) {
return new JdbcInstanceTokenRepository(connectionProvider, ordoSqlDialect);
}
@Bean
@ConditionalOnMissingBean
public OrdoEngine ordoEngine(Clock ordoClock,
@@ -146,11 +155,12 @@ public class OrdoJdbcAutoConfiguration {
ApprovalTaskRepository ordoApprovalTaskRepository,
ProcessHistoryRepository ordoProcessHistoryRepository,
ActionExecutionRepository ordoActionExecutionRepository,
InstanceTokenRepository ordoInstanceTokenRepository,
ObjectProvider<OrdoEventListener> ordoEventListeners) {
return new DefaultOrdoEngine(ordoClock, ordoAssigneeResolver, ordoRoutingCondition, actionHandler,
ordoTransactionExecutor, ordoProcessDefinitionRepository, ordoProcessInstanceRepository,
ordoApprovalTaskRepository, ordoProcessHistoryRepository, ordoActionExecutionRepository,
ordoEventListeners.orderedStream().toList());
ordoInstanceTokenRepository, ordoEventListeners.orderedStream().toList());
}
@Bean
@@ -0,0 +1,164 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.InstanceToken;
import com.jetlumen.ordo.api.TokenStatus;
import com.jetlumen.ordo.api.repository.InstanceTokenRepository;
import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect;
import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
/** JDBC implementation of the parallel token port. */
public final class JdbcInstanceTokenRepository implements InstanceTokenRepository {
private static final String COLUMNS =
"id, instance_id, parallel_step_id, branch_id, current_step_id, status";
private static final String INSERT =
"INSERT INTO ordo_instance_token (id, instance_id, parallel_step_id, branch_id, current_step_id, status)"
+ " VALUES (?, ?, ?, ?, ?, ?)";
private static final String SELECT_BY_INSTANCE =
"SELECT " + COLUMNS + " FROM ordo_instance_token WHERE instance_id = ? ORDER BY parallel_step_id, branch_id";
private static final String SELECT_BY_PARALLEL =
"SELECT " + COLUMNS + " FROM ordo_instance_token WHERE instance_id = ? AND parallel_step_id = ?"
+ " ORDER BY branch_id";
private static final String COMPLETE_IF_ACTIVE =
"UPDATE ordo_instance_token SET status = 'COMPLETED' WHERE instance_id = ? AND parallel_step_id = ?"
+ " AND branch_id = ? AND status = 'ACTIVE'";
private static final String UPDATE_CURRENT_IF_ACTIVE =
"UPDATE ordo_instance_token SET current_step_id = ? WHERE instance_id = ? AND parallel_step_id = ?"
+ " AND branch_id = ? AND status = 'ACTIVE'";
private static final String COMPLETE_ACTIVE_BY_INSTANCE =
"UPDATE ordo_instance_token SET status = 'COMPLETED' WHERE instance_id = ? AND status = 'ACTIVE'";
private final JdbcConnectionProvider connectionProvider;
public JdbcInstanceTokenRepository(JdbcConnectionProvider connectionProvider) {
this(connectionProvider, SqlDialects.postgresql());
}
public JdbcInstanceTokenRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
Objects.requireNonNull(dialect, "dialect must not be null");
}
@Override
public void insert(InstanceToken token) {
Objects.requireNonNull(token, "token must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement insert = connection.prepareStatement(INSERT)) {
insert.setString(1, token.id());
insert.setString(2, token.instanceId());
insert.setString(3, token.parallelStepId());
insert.setString(4, token.branchId());
insert.setString(5, token.currentStepId());
insert.setString(6, token.status().name());
insert.executeUpdate();
} catch (SQLException e) {
throw new JdbcStorageException("failed to insert token: " + token.id(), e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public List<InstanceToken> findByInstanceId(String instanceId) {
Objects.requireNonNull(instanceId, "instanceId must not be null");
return query(SELECT_BY_INSTANCE, instanceId);
}
@Override
public List<InstanceToken> findByInstanceAndParallel(String instanceId, String parallelStepId) {
Objects.requireNonNull(instanceId, "instanceId must not be null");
Objects.requireNonNull(parallelStepId, "parallelStepId must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(SELECT_BY_PARALLEL)) {
select.setString(1, instanceId);
select.setString(2, parallelStepId);
return readAll(select);
} catch (SQLException e) {
throw new JdbcStorageException("failed to load tokens for parallel: " + parallelStepId, e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public boolean completeIfActive(String instanceId, String parallelStepId, String branchId) {
return update(COMPLETE_IF_ACTIVE, instanceId, parallelStepId, branchId);
}
@Override
public boolean updateCurrentStepIfActive(String instanceId, String parallelStepId, String branchId,
String currentStepId) {
Objects.requireNonNull(currentStepId, "currentStepId must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(UPDATE_CURRENT_IF_ACTIVE)) {
update.setString(1, currentStepId);
update.setString(2, instanceId);
update.setString(3, parallelStepId);
update.setString(4, branchId);
return update.executeUpdate() == 1;
} catch (SQLException e) {
throw new JdbcStorageException("failed to update token current step", e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public void completeActiveByInstanceId(String instanceId) {
Objects.requireNonNull(instanceId, "instanceId must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(COMPLETE_ACTIVE_BY_INSTANCE)) {
update.setString(1, instanceId);
update.executeUpdate();
} catch (SQLException e) {
throw new JdbcStorageException("failed to complete tokens for instance: " + instanceId, e);
} finally {
connectionProvider.close(connection);
}
}
private List<InstanceToken> query(String sql, String instanceId) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) {
select.setString(1, instanceId);
return readAll(select);
} catch (SQLException e) {
throw new JdbcStorageException("failed to load tokens for instance: " + instanceId, e);
} finally {
connectionProvider.close(connection);
}
}
private boolean update(String sql, String instanceId, String parallelStepId, String branchId) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(sql)) {
update.setString(1, instanceId);
update.setString(2, parallelStepId);
update.setString(3, branchId);
return update.executeUpdate() == 1;
} catch (SQLException e) {
throw new JdbcStorageException("failed to update token", e);
} finally {
connectionProvider.close(connection);
}
}
private static List<InstanceToken> readAll(PreparedStatement select) throws SQLException {
List<InstanceToken> tokens = new ArrayList<>();
try (ResultSet resultSet = select.executeQuery()) {
while (resultSet.next()) {
tokens.add(new InstanceToken(resultSet.getString("id"), resultSet.getString("instance_id"),
resultSet.getString("parallel_step_id"), resultSet.getString("branch_id"),
resultSet.getString("current_step_id"), TokenStatus.valueOf(resultSet.getString("status"))));
}
}
return tokens;
}
}
@@ -1,7 +1,9 @@
package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.StepKind;
import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.query.Page;
import com.jetlumen.ordo.api.query.PageRequest;
@@ -19,8 +21,10 @@ import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
import java.sql.Types;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -39,8 +43,8 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
"INSERT INTO ordo_process_definition (id, version, name, created_at) VALUES (?, ?, ?, ?)";
private static final String INSERT_STEP =
"INSERT INTO ordo_approval_step (definition_id, definition_version, step_id, step_name, policy, step_order,"
+ " kind, action_key, due_after, due_then, due_to, due_action)"
+ " VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
+ " kind, action_key, due_after, due_then, due_to, due_action, parent_step_id, branch_id, branch_order)"
+ " VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
private static final String INSERT_CANDIDATE =
"INSERT INTO ordo_step_candidate (definition_id, definition_version, step_id, candidate, candidate_order)"
+ " VALUES (?, ?, ?, ?, ?)";
@@ -50,7 +54,8 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
private static final String SELECT_DEFINITION =
"SELECT name FROM ordo_process_definition WHERE id = ? AND version = ?";
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,"
+ " parent_step_id, branch_id, branch_order"
+ " FROM ordo_approval_step"
+ " WHERE definition_id = ? AND definition_version = ? ORDER BY step_order";
private static final String SELECT_CANDIDATES =
@@ -248,12 +253,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
}
}
}
List<ApprovalStep> steps = new ArrayList<>();
for (StepRow row : stepRows) {
List<String> candidates = candidatesByStep.getOrDefault(row.stepId(), List.of());
steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(),
row.actionKey(), ApprovalStepMapper.toDue(row)));
}
List<ApprovalStep> steps = assembleSteps(stepRows, candidatesByStep);
List<StepTransition> transitions = new ArrayList<>();
try (PreparedStatement selectTransitions = connection.prepareStatement(SELECT_TRANSITIONS)) {
selectTransitions.setString(1, definitionId);
@@ -269,6 +269,41 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
return new ProcessDefinition(definitionId, version, name, steps, transitions);
}
private static List<ApprovalStep> assembleSteps(List<StepRow> stepRows, Map<String, List<String>> candidatesByStep) {
Map<String, LinkedHashMap<String, List<String>>> branchesByParent = new LinkedHashMap<>();
Map<String, Integer> branchOrderByKey = new LinkedHashMap<>();
for (StepRow row : stepRows) {
if (row.parentStepId() == null || row.parentStepId().isBlank()) {
continue;
}
branchOrderByKey.putIfAbsent(row.parentStepId() + '\0' + row.branchId(),
row.branchOrder() == null ? Integer.MAX_VALUE : row.branchOrder());
branchesByParent.computeIfAbsent(row.parentStepId(), key -> new LinkedHashMap<>())
.computeIfAbsent(row.branchId(), key -> new ArrayList<>())
.add(row.stepId());
}
List<ApprovalStep> steps = new ArrayList<>();
for (StepRow row : stepRows) {
List<String> candidates = candidatesByStep.getOrDefault(row.stepId(), List.of());
List<ParallelBranch> branches = List.of();
if (row.kind() == StepKind.PARALLEL) {
LinkedHashMap<String, List<String>> byBranch = branchesByParent.getOrDefault(row.stepId(),
new LinkedHashMap<>());
List<Map.Entry<String, List<String>>> ordered = new ArrayList<>(byBranch.entrySet());
ordered.sort(Comparator.comparingInt(entry ->
branchOrderByKey.getOrDefault(row.stepId() + '\0' + entry.getKey(), Integer.MAX_VALUE)));
List<ParallelBranch> assembled = new ArrayList<>();
for (Map.Entry<String, List<String>> entry : ordered) {
assembled.add(new ParallelBranch(entry.getKey(), entry.getValue()));
}
branches = assembled;
}
steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(),
row.actionKey(), ApprovalStepMapper.toDue(row), branches));
}
return steps;
}
private Integer lockCurrentVersion(Connection connection, String definitionId) throws SQLException {
try (PreparedStatement select = connection.prepareStatement(lockProcess)) {
select.setString(1, definitionId);
@@ -336,6 +371,23 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
}
private static void insertGraph(Connection connection, ProcessDefinition definition) throws SQLException {
Map<String, String> parentByStep = new LinkedHashMap<>();
Map<String, String> branchByStep = new LinkedHashMap<>();
Map<String, Integer> branchOrderByStep = new LinkedHashMap<>();
for (ApprovalStep step : definition.steps()) {
if (step.kind() != StepKind.PARALLEL) {
continue;
}
int branchOrder = 0;
for (ParallelBranch branch : step.branches()) {
for (String nestedId : branch.stepIds()) {
parentByStep.put(nestedId, step.id());
branchByStep.put(nestedId, branch.id());
branchOrderByStep.put(nestedId, branchOrder);
}
branchOrder++;
}
}
int stepOrder = 0;
for (ApprovalStep step : definition.steps()) {
try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) {
@@ -358,6 +410,14 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
SqlBindings.setString(insertStep, 11, step.due().to());
SqlBindings.setString(insertStep, 12, step.due().action());
}
SqlBindings.setString(insertStep, 13, parentByStep.get(step.id()));
SqlBindings.setString(insertStep, 14, branchByStep.get(step.id()));
Integer branchOrder = branchOrderByStep.get(step.id());
if (branchOrder == null) {
insertStep.setNull(15, Types.INTEGER);
} else {
insertStep.setInt(15, branchOrder);
}
insertStep.executeUpdate();
}
int candidateOrder = 0;
@@ -26,7 +26,15 @@ public final class ApprovalStepMapper {
resultSet.getString("due_after"),
resultSet.getString("due_then"),
resultSet.getString("due_to"),
resultSet.getString("due_action"));
resultSet.getString("due_action"),
resultSet.getString("parent_step_id"),
resultSet.getString("branch_id"),
intOrNull(resultSet, "branch_order"));
}
private static Integer intOrNull(ResultSet resultSet, String column) throws SQLException {
int value = resultSet.getInt(column);
return resultSet.wasNull() ? null : value;
}
public static StepDue toDue(StepRow row) {
@@ -38,6 +46,7 @@ public final class ApprovalStepMapper {
}
public record StepRow(String stepId, String stepName, ApprovalPolicy policy, StepKind kind, String actionKey,
String dueAfter, String dueThen, String dueTo, String dueAction) {
String dueAfter, String dueThen, String dueTo, String dueAction, String parentStepId,
String branchId, Integer branchOrder) {
}
}
@@ -0,0 +1,20 @@
-- Structured PARALLEL blocks: parent/branch on steps, runtime tokens.
ALTER TABLE ordo_approval_step ADD COLUMN parent_step_id VARCHAR(64);
ALTER TABLE ordo_approval_step ADD COLUMN branch_id VARCHAR(64);
ALTER TABLE ordo_approval_step ADD COLUMN branch_order INTEGER;
CREATE TABLE ordo_instance_token (
id VARCHAR(36) NOT NULL,
instance_id VARCHAR(36) NOT NULL,
parallel_step_id VARCHAR(64) NOT NULL,
branch_id VARCHAR(64) NOT NULL,
current_step_id VARCHAR(64) NOT NULL,
status VARCHAR(16) NOT NULL,
PRIMARY KEY (id),
CONSTRAINT fk_instance_token_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id),
CONSTRAINT uq_instance_token_branch UNIQUE (instance_id, parallel_step_id, branch_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
CREATE INDEX idx_instance_token_instance ON ordo_instance_token (instance_id);
CREATE INDEX idx_instance_token_parallel ON ordo_instance_token (instance_id, parallel_step_id);
@@ -0,0 +1,19 @@
-- Structured PARALLEL blocks: parent/branch on steps, runtime tokens.
ALTER TABLE ordo_approval_step ADD COLUMN parent_step_id VARCHAR(64);
ALTER TABLE ordo_approval_step ADD COLUMN branch_id VARCHAR(64);
ALTER TABLE ordo_approval_step ADD COLUMN branch_order INTEGER;
CREATE TABLE ordo_instance_token (
id VARCHAR(36) PRIMARY KEY,
instance_id VARCHAR(36) NOT NULL,
parallel_step_id VARCHAR(64) NOT NULL,
branch_id VARCHAR(64) NOT NULL,
current_step_id VARCHAR(64) NOT NULL,
status VARCHAR(16) NOT NULL,
CONSTRAINT fk_instance_token_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id),
CONSTRAINT uq_instance_token_branch UNIQUE (instance_id, parallel_step_id, branch_id)
);
CREATE INDEX idx_instance_token_instance ON ordo_instance_token (instance_id);
CREATE INDEX idx_instance_token_parallel ON ordo_instance_token (instance_id, parallel_step_id);
@@ -5,11 +5,13 @@ import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.AssigneeResolver;
import com.jetlumen.ordo.api.OrdoEngine;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus;
import com.jetlumen.ordo.api.RoutingCondition;
import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.TaskAction;
import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.core.DefaultOrdoEngine;
@@ -76,7 +78,8 @@ class JdbcMysqlIntegrationTest {
}
MysqlDataSource schemaDataSource = newDataSource(DATABASE);
JdbcTestSupport.applySchema(schemaDataSource, JdbcTestSupport.MYSQL_BASELINE, JdbcTestSupport.MYSQL_V2);
JdbcTestSupport.applySchema(schemaDataSource, JdbcTestSupport.MYSQL_BASELINE, JdbcTestSupport.MYSQL_V2,
JdbcTestSupport.MYSQL_V3);
dataSource = schemaDataSource;
}
@@ -118,6 +121,19 @@ class JdbcMysqlIntegrationTest {
assertEquals("LEAVE-2026-001", finished.context().value("requestId").orElseThrow());
}
@Test
void engineJoinsAParallelBlock() {
OrdoEngine engine = newEngine(AssigneeResolver.direct());
engine.publish(parallelReview("parallel-mysql"));
ProcessInstance instance = engine.start("parallel-mysql", "alice");
assertEquals(2, engine.findPendingTasksByInstanceId(instance.id()).size());
for (ApprovalTask task : engine.findPendingTasksByInstanceId(instance.id())) {
engine.approve(task.id(), task.assignee());
}
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "cara");
assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void publishIsIdempotentForTheSameGraphAndVersionsAChangedGraph() {
JdbcProcessDefinitionRepository repository =
@@ -233,9 +249,25 @@ class JdbcMysqlIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider, DIALECT),
new JdbcProcessHistoryRepository(connectionProvider, DIALECT),
new JdbcActionExecutionRepository(connectionProvider, DIALECT),
new JdbcInstanceTokenRepository(connectionProvider, DIALECT),
List.of());
}
private static ProcessDefinition parallelReview(String id) {
return new ProcessDefinition(id, "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo")));
}
private static MysqlDataSource newDataSource(String database) {
MysqlDataSource mysqlDataSource = new MysqlDataSource();
String host = JdbcTestSupport.config("ordo.test.mysql.host", "ORDO_TEST_MYSQL_HOST", "localhost");
@@ -7,6 +7,7 @@ import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.AssigneeResolver;
import com.jetlumen.ordo.api.OrdoEngine;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessEvent;
@@ -134,6 +135,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
startEngine.publish(dueDefinition);
ProcessInstance instance = startEngine.start("leave-due", "alice");
@@ -147,6 +149,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
assertEquals(1, laterEngine.processDue(10));
ApprovalTask escalated = laterEngine.findPendingTasksByInstanceId(instance.id()).get(0);
@@ -190,6 +193,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
failingEngine.publish(new ProcessDefinition("leave-noroute", "Leave request",
List.of(ApprovalStep.single("manager", "Manager approval", "maria")),
@@ -298,6 +302,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
actionEngine.publish(new ProcessDefinition("leave-action", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"),
@@ -323,6 +328,29 @@ class JdbcOrdoEngineIntegrationTest {
assertTrue(types.contains(ProcessEventType.INSTANCE_APPROVED));
}
@Test
void joinsAParallelBlock() {
engine.publish(new ProcessDefinition("parallel-review", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo"))));
ProcessInstance instance = engine.start("parallel-review", "alice");
assertEquals(2, engine.findPendingTasksByInstanceId(instance.id()).size());
for (ApprovalTask task : List.copyOf(engine.findPendingTasksByInstanceId(instance.id()))) {
engine.approve(task.id(), task.assignee());
}
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "cara");
assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status());
}
private OrdoEngine newEngine(AssigneeResolver assigneeResolver) {
return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver, RoutingCondition.always(),
ActionHandler.noop(),
@@ -332,6 +360,7 @@ class JdbcOrdoEngineIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
}
}
@@ -5,11 +5,13 @@ import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.AssigneeResolver;
import com.jetlumen.ordo.api.OrdoEngine;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus;
import com.jetlumen.ordo.api.RoutingCondition;
import com.jetlumen.ordo.api.StepTransition;
import com.jetlumen.ordo.api.TaskAction;
import com.jetlumen.ordo.api.TaskStatus;
import com.jetlumen.ordo.core.DefaultOrdoEngine;
@@ -121,6 +123,19 @@ class JdbcPostgresIntegrationTest {
assertEquals("LEAVE-2026-001", finished.context().value("requestId").orElseThrow());
}
@Test
void engineJoinsAParallelBlock() {
OrdoEngine engine = newEngine(AssigneeResolver.direct());
engine.publish(parallelReview("parallel-pg"));
ProcessInstance instance = engine.start("parallel-pg", "alice");
assertEquals(2, engine.findPendingTasksByInstanceId(instance.id()).size());
for (ApprovalTask task : engine.findPendingTasksByInstanceId(instance.id())) {
engine.approve(task.id(), task.assignee());
}
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "cara");
assertEquals(ProcessStatus.APPROVED, engine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void publishIsIdempotentForTheSameGraphAndVersionsAChangedGraph() {
JdbcProcessDefinitionRepository repository = new JdbcProcessDefinitionRepository(connectionProvider);
@@ -229,9 +244,25 @@ class JdbcPostgresIntegrationTest {
new JdbcApprovalTaskRepository(connectionProvider),
new JdbcProcessHistoryRepository(connectionProvider),
new JdbcActionExecutionRepository(connectionProvider),
new JdbcInstanceTokenRepository(connectionProvider),
List.of());
}
private static ProcessDefinition parallelReview(String id) {
return new ProcessDefinition(id, "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo")));
}
private static PGSimpleDataSource newDataSource(String currentSchema) {
PGSimpleDataSource pgDataSource = new PGSimpleDataSource();
pgDataSource.setServerNames(new String[]{JdbcTestSupport.config("ordo.test.pg.host", "ORDO_TEST_PG_HOST", "localhost")});
@@ -2,6 +2,7 @@ package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalPolicy;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ParallelBranch;
import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.StepDue;
import com.jetlumen.ordo.api.StepKind;
@@ -97,6 +98,24 @@ class JdbcProcessDefinitionRepositoryTest {
assertEquals("leave-approved-mail", repository.findLatest("notify").orElseThrow().steps().get(1).actionKey());
}
@Test
void insertsAndReadsBackAParallelBlock() {
ProcessDefinition definition = new ProcessDefinition("parallel-review", "P", List.of(
ApprovalStep.parallel("dept", "Dept", List.of(
new ParallelBranch("legal", List.of("legal")),
new ParallelBranch("finance", List.of("finance")))),
ApprovalStep.single("legal", "Legal", "lee"),
ApprovalStep.single("finance", "Finance", "fay"),
ApprovalStep.single("ceo", "CEO", "cara")),
List.of(
StepTransition.end("legal"),
StepTransition.end("finance"),
StepTransition.always("dept", "ceo"),
StepTransition.end("ceo")));
repository.publish(definition);
assertEquals(definition.withVersion(1), repository.findLatest("parallel-review").orElseThrow());
}
private static ProcessDefinition definition(String id, String name) {
return ProcessDefinition.linear(id, name, List.of(ApprovalStep.single("lead", "Lead approval", "lee")));
}
@@ -15,8 +15,10 @@ import java.util.UUID;
final class JdbcTestSupport {
static final String POSTGRES_BASELINE = "/db/postgresql/migration/V1__baseline.sql";
static final String POSTGRES_V2 = "/db/postgresql/migration/V2__transition_condition_json.sql";
static final String POSTGRES_V3 = "/db/postgresql/migration/V3__parallel_blocks.sql";
static final String MYSQL_BASELINE = "/db/mysql/migration/V1__baseline.sql";
static final String MYSQL_V2 = "/db/mysql/migration/V2__transition_condition_json.sql";
static final String MYSQL_V3 = "/db/mysql/migration/V3__parallel_blocks.sql";
private JdbcTestSupport() {
}
@@ -43,7 +45,7 @@ final class JdbcTestSupport {
}
static void applySchema(DataSource dataSource) {
applySchema(dataSource, POSTGRES_BASELINE, POSTGRES_V2);
applySchema(dataSource, POSTGRES_BASELINE, POSTGRES_V2, POSTGRES_V3);
}
static void applySchema(DataSource dataSource, String... resourcePaths) {