feat: let the current assignee reassign a pending task

Keep the same task id, move the pending inbox, and record TASK_REASSIGNED without advancing the step.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-16 08:40:25 +08:00
co-authored by Cursor
parent 5795778112
commit 850583f326
15 changed files with 286 additions and 26 deletions
+7 -5
View File
@@ -1,6 +1,6 @@
# Ordo
轻量审批流程引擎。宿主通过 `OrdoEngine` 注册流程定义、发起实例、审批/驳回/撤回,并查询任务、实例与审计历史。引擎不绑定业务表单,也不自带 REST;业务数据放在 `ProcessContext` 里。
轻量审批流程引擎。宿主通过 `OrdoEngine` 注册流程定义、发起实例、审批/驳回/转派/撤回,并查询任务、实例与审计历史。引擎不绑定业务表单、用户体系或设计器 UI,当前也不自带 REST;业务数据放在 `ProcessContext` 里。
要求 **Java 17+**。当前版本 `0.0.1-SNAPSHOT`。
@@ -22,11 +22,12 @@
- 会签/或签:`ApprovalPolicy.ALL` / `ANY`(多候选人)
- ACTION 步骤:事务提交后调用宿主 `ActionHandler`
- 发起人撤回:`WITHDRAWN`,待办任务 `SKIPPED`
- 任务转派:当前办理人 `reassign`,审计 `TASK_REASSIGNED`
- 分页查询:任务 / 实例 / 流程定义
- 审计时间线:`ProcessEvent` + `queryHistory`
- 扩展点:`AssigneeResolver`、`RoutingCondition`、`ActionHandler`、`OrdoEventListener`
开发计划:任务转派、到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言。多租户与官方 REST Starter **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。
开发计划:到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST、目录与多版本定义之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。
详细用法(定义 JSON、扩展点、异常、查询、ACTION/审计语义)见 **[docs/usage.md](docs/usage.md)**。对外行为变更时同步更新该文档。
@@ -106,13 +107,14 @@ ordo:
| Bean | 默认 |
|---|---|
| `AssigneeResolver` | 候选人即办理人 |
| `RoutingCondition` | 始终匹配 |
| `RoutingCondition` | 始终匹配(全局单例;宿主可按 key 分发) |
| `ActionHandler` | 空操作 |
| `OrdoEventListener` | 可注册多个,提交后按顺序调用 |
## 运行时约定
**办理人** 必须等于任务 `assignee`,否则 `UnauthorizedTaskOperationException`。
**办理人** 必须等于任务 `assignee`,否则 `UnauthorizedTaskOperationException`(`approve` / `reject` / `reassign`)。
**转派** 只改 PENDING 任务的 `assignee`,不推进步骤。
**撤回** 仅发起人可操作,且实例须为 `RUNNING`。
实例状态:`RUNNING` / `APPROVED` / `REJECTED` / `WITHDRAWN`。
@@ -120,7 +122,7 @@ ordo:
**ACTION** 在审批事务提交之后执行。失败只记 `FAILED`、打日志、发 `ACTION_FAILED`,**不回滚已生效审批、不阻塞后续步骤**。可重试策略留给宿主(listener 或 `queryActionExecutions`)。
**历史** `queryHistory(instanceId, page)` 按时间升序。事件类型包括实例起止、任务创建/审批/跳过、ACTION 成败。
**历史** `queryHistory(instanceId, page)` 按时间升序。事件类型包括实例起止、任务创建/审批/跳过/转派、ACTION 成败。
`OrdoEventListener.onEvent` 在提交后派发;单个 listener 抛错不影响流程和其他 listener。
## 查询
+17 -6
View File
@@ -13,15 +13,12 @@
- 发起人撤回 `WITHDRAWN`
- 分页查询:任务 / 实例 / 定义
- 审计 `ProcessEvent` / `queryHistory`;`OrdoEventListener`
- 任务转派 `reassign` / `TASK_REASSIGNED`
## 开发计划(确定要做)
下列能力已纳入计划,尚未实现。实现顺序可按依赖调整,但范围本身不从计划中拿掉。
### 任务转派
任务创建后当前 `assignee` 不可变。计划提供委托/转派(delegate/reassign),把待办转给他人办理,并写入审计。
### 到期升级
目前无到期时间与定时器。计划支持 SLA/超时:到期后升级(改办理人、通知或进入指定步骤),即「到期升级」。
@@ -38,13 +35,27 @@
当前 JDBC DDL/冲突处理面向 PostgreSQL。计划增加 MySQL 方言(及对应测试),通过 dialect 层扩展,而不是只支持一种库。
### 可选 REST + 目录 SPI
引擎入口仍是 `OrdoEngine`。计划提供**可选、极薄**的 REST 适配(例如独立 starter),带 OpenAPI,覆盖定义读写/解析校验、实例与任务查询,不包含鉴权、RBAC、业务表单。
配套 **目录 SPI**:宿主登记可用的 `conditionKey` / `actionKey` / 候选人(或角色)项,供 REST 与外部设计器下拉,而不是在 JSON 里写引擎无法执行的表达式。
`RoutingCondition` / `ActionHandler` / `AssigneeResolver` 保持全局单例;流程隔离由宿主用 key 约定(建议前缀)+ 门面分发,引擎不按流程定义拆 bean。
### 独立设计器(不进本仓库)
流程设计器是**单独产品**,消费上述 REST 与目录,不做成 ordo 模块。画布对齐引擎图(审批步、ACTION 步、边上的 `when`/`priority`,结束为 `to: null`),不引入 BPMN 网关/并行等引擎没有的语义。节点坐标等 layout 由设计器自存,不进入 `ProcessDefinition`。
启动时机:REST 契约、目录 SPI、定义不可变多版本落地之后。当前 `replace` + 无版本锁定不适合作为设计器保存/发布模型。
## 暂不在计划中
- **多租户**(数据模型无 tenant 隔离)
- **官方 REST Starter**(不提供 `ordo-spring-boot-starter-web`;REST 由宿主自建)
- **设计器 UI**(不进 ordo;见上节独立产品)
## 后续新特性
上表「确定要做」之外,仍可能立项其他能力(例如表单/UI schema、子流程、并行 fork-join)。**多租户与官方 REST 在另有明确决定前不进入计划。**
上表「确定要做」之外,仍可能立项其他能力(例如表单/UI schema、子流程、并行 fork-join)。**多租户在另有明确决定前不进入计划。**
新特性立项时写入「开发计划」对应小节;完成后移到「已完成」,并更新 [usage.md](usage.md)。
+16 -11
View File
@@ -10,9 +10,9 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。
做:流程定义、实例推进、待办任务、条件路由、ACTION 副作用、审计事件、分页查询。
产品边界(不做,且暂不在开发计划):业务表单、用户体系、官方 REST、多租户。业务字段放在 `ProcessContext`(不可变 `Map<String, Object>`)。
产品边界:不做业务表单、用户体系、多租户;不内置设计器 UI。业务字段放在 `ProcessContext`(不可变 `Map<String, Object>`)。当前也**没有** REST;HTTP 仍由宿主自建。计划中的可选 REST 与独立设计器见 [roadmap.md](roadmap.md)。
开发计划(尚未提供,见 [roadmap.md](roadmap.md)):任务转派、到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言。
开发计划(尚未提供,见 [roadmap.md](roadmap.md)):到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。
## 2. 模块与接入
@@ -189,7 +189,8 @@ ProcessInstance instance = ordo.start("leave-request", "alice",
new ProcessContext(Map.of("requestId", "LEAVE-001", "days", 2)));
List<ApprovalTask> pending = ordo.findPendingTasksByInstanceId(instance.id());
ordo.approve(pending.get(0).id(), "maria", "ok");
ordo.reassign(pending.get(0).id(), "maria", "diana");
ordo.approve(pending.get(0).id(), "diana", "ok");
ordo.reject(taskId, "maria", "额度不足");
ordo.withdraw(instance.id(), "alice", "计划有变");
@@ -199,11 +200,12 @@ ordo.withdraw(instance.id(), "alice", "计划有变");
### 4.2 权限
- `approve` / `reject`:`actor` 必须等于该任务当前 `assignee`,否则 `UnauthorizedTaskOperationException`。
- `approve` / `reject` / `reassign`:`actor` 必须等于该任务当前 `assignee`,否则 `UnauthorizedTaskOperationException`。
- `reassign`:仅 `PENDING` 任务;同一任务 id,办理人改为 `newAssignee`,不推进步骤。`newAssignee` 不可空白、不可等于当前 `assignee`,且同一步不能已有该人的 `PENDING` 任务,否则 `IllegalArgumentException`。不经过 `AssigneeResolver`。
- 任务非 `PENDING`:`TaskAlreadyCompletedException`。
- `withdraw`:仅 `initiator`,否则 `UnauthorizedInstanceOperationException`;实例非 `RUNNING`:`InstanceAlreadyCompletedException`。
没有转派、没有管理员代批、没有系统取消。
没有管理员代批、没有系统取消。
### 4.3 状态
@@ -237,7 +239,9 @@ ordo.withdraw(instance.id(), "alice", "计划有变");
无条件边通常作为默认分支,`priority` 应大于带 `when` 的边。
`RoutingCondition` 只看到 `ProcessContext`,看不到任务意见。上下文在 `start` 时写入,运行中引擎**不会**改 context。
`RoutingCondition` 只看到 `ProcessContext`,看不到任务意见,也**没有**流程定义 id。上下文在 `start` 时写入,运行中引擎**不会**改 context。
引擎只注入**一个** `RoutingCondition`(与 `ActionHandler` 相同)。Spring 下多个该类型 Bean 会冲突。宿主用一个门面按 `conditionKey` 分发到多套规则;不同流程靠 key 约定隔离(例如 `leave.days-gt-3`),不要指望引擎按定义拆 bean。
## 7. ACTION 步骤
@@ -271,10 +275,10 @@ public class MailActions implements ActionHandler {
`ProcessEventType`:
- `INSTANCE_STARTED` / `INSTANCE_APPROVED` / `INSTANCE_REJECTED` / `INSTANCE_WITHDRAWN`
- `TASK_CREATED` / `TASK_APPROVED` / `TASK_REJECTED` / `TASK_SKIPPED`
- `TASK_CREATED` / `TASK_APPROVED` / `TASK_REJECTED` / `TASK_SKIPPED` / `TASK_REASSIGNED`
- `ACTION_SUCCEEDED` / `ACTION_FAILED`
字段:`id`、`instanceId`、可选 `taskId`/`stepId`/`actor`/`detail`、`occurredAt`。系统完成类事件 `actor` 可为空。
字段:`id`、`instanceId`、可选 `taskId`/`stepId`/`actor`/`detail`、`occurredAt`。系统完成类事件 `actor` 可为空。`TASK_REASSIGNED` 的 `actor` 为转出人,`detail` 为转入人。
`OrdoEventListener.onEvent(ProcessEvent)` 在**事务提交之后**按事件顺序调用(含本轮 ACTION 结果)。单个 listener 抛错只打日志,不影响流程和其他 listener。
@@ -307,6 +311,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| `register` / `replace` | 登记 / 整图替换定义 |
| `start` | 发起;可选 `ProcessContext` |
| `approve` / `reject` | 办理当前 PENDING 任务 |
| `reassign` | 当前办理人把 PENDING 任务转给他人 |
| `withdraw` | 发起人撤回 |
| `find*` | 按 id / 待办索引读取 |
| `queryTasks` / `queryInstances` / `queryDefinitions` | 分页列表 |
@@ -325,7 +330,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| `InstanceAlreadyCompletedException` | 对非 RUNNING 实例完成/撤回 |
| `TaskNotFoundException` | 任务 id 不存在 |
| `TaskAlreadyCompletedException` | 重复办理或已被 SKIPPED |
| `UnauthorizedTaskOperationException` | actor ≠ assignee |
| `UnauthorizedTaskOperationException` | actor ≠ 当前 assignee(approve / reject / reassign) |
| `UnauthorizedInstanceOperationException` | 非发起人撤回 |
| `NoRouteFoundException` | 当前步没有匹配转移 |
@@ -339,6 +344,6 @@ Flyway 脚本在 `ordo-storage-jdbc` 的 `db/migration`(V1–V5)。表包括
## 13. 未提供能力
开发计划中(见 [roadmap.md](roadmap.md)):转派、到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言。
开发计划中(见 [roadmap.md](roadmap.md)):到期升级、`CANCELLED`、定义不可变多版本、MySQL 方言、可选 REST + 目录 SPI。独立设计器不进本仓库,等 REST、目录与多版本定义之后再做。
暂不在计划中:多租户、官方 REST Starter。REST 由宿主自建。
暂不在计划中:多租户、设计器 UI。当前 REST 由宿主自建。
@@ -29,6 +29,7 @@ public interface OrdoEngine {
return reject(taskId, actor, null);
}
ApprovalTask reject(String taskId, String actor, String comment);
ApprovalTask reassign(String taskId, String actor, String newAssignee);
default ProcessInstance withdraw(String instanceId, String actor) {
return withdraw(instanceId, actor, null);
}
@@ -7,6 +7,7 @@ public enum ProcessEventType {
TASK_APPROVED,
TASK_REJECTED,
TASK_SKIPPED,
TASK_REASSIGNED,
INSTANCE_APPROVED,
INSTANCE_REJECTED,
INSTANCE_WITHDRAWN,
@@ -25,12 +25,21 @@ public interface ApprovalTaskRepository {
List<ApprovalTask> findPendingByInstanceId(String instanceId);
/**
* Atomically completes the task only if it is still pending.
* Atomically completes the task only if it is still pending and still assigned to
* {@code completedTask.assignee()}.
*
* @return true if the update was applied, false if the task had already been completed by a concurrent operation
* @return true if the update was applied, false if the task had already been completed or reassigned
*/
boolean completeIfPending(ApprovalTask completedTask);
/**
* Atomically changes the assignee only if the task is still pending and still assigned to
* {@code expectedAssignee}.
*
* @return true if the update was applied
*/
boolean reassignIfPending(String taskId, String expectedAssignee, String newAssignee);
/** Paginated, filterable query; results are ordered newest-first (created_at desc). */
Page<ApprovalTask> query(TaskQuery query, PageRequest pageRequest);
}
@@ -174,6 +174,36 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return completed;
}
@Override
public synchronized ApprovalTask reassign(String taskId, String actor, String newAssignee) {
requireText(newAssignee, "new assignee");
List<ProcessEvent> events = new ArrayList<>();
ApprovalTask reassigned = transactionExecutor.execute(() -> {
ApprovalTask task = requirePendingTaskForActor(taskId, actor);
if (task.assignee().equals(newAssignee)) {
throw new IllegalArgumentException("new assignee must differ from current assignee");
}
boolean duplicatePending = taskRepository.findByInstanceIdAndStepId(task.instanceId(), task.stepId()).stream()
.anyMatch(other -> !other.id().equals(task.id())
&& other.status() == TaskStatus.PENDING
&& other.assignee().equals(newAssignee));
if (duplicatePending) {
throw new IllegalArgumentException("new assignee already has a pending task on this step");
}
if (!taskRepository.reassignIfPending(task.id(), task.assignee(), newAssignee)) {
throw new TaskAlreadyCompletedException(task.id());
}
Instant now = clock.instant();
ApprovalTask updated = new ApprovalTask(task.id(), task.instanceId(), task.stepId(), task.name(),
newAssignee, task.status(), task.createdAt(), task.completedAt(), task.action());
record(events, task.instanceId(), task.id(), task.stepId(), ProcessEventType.TASK_REASSIGNED, actor,
newAssignee, now);
return updated;
});
dispatch(events);
return reassigned;
}
@Override
public synchronized ProcessInstance withdraw(String instanceId, String actor, String comment) {
requireText(instanceId, "instance id");
@@ -100,6 +100,11 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
return delegate.reject(taskId, actor, comment);
}
@Override
public ApprovalTask reassign(String taskId, String actor, String newAssignee) {
return delegate.reassign(taskId, actor, newAssignee);
}
@Override
public ProcessInstance withdraw(String instanceId, String actor, String comment) {
return delegate.withdraw(instanceId, actor, comment);
@@ -76,13 +76,26 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
@Override
public synchronized boolean completeIfPending(ApprovalTask completedTask) {
ApprovalTask current = tasks.get(completedTask.id());
if (current == null || current.status() != TaskStatus.PENDING) {
if (current == null || current.status() != TaskStatus.PENDING
|| !current.assignee().equals(completedTask.assignee())) {
return false;
}
tasks.put(completedTask.id(), completedTask);
return true;
}
@Override
public synchronized boolean reassignIfPending(String taskId, String expectedAssignee, String newAssignee) {
ApprovalTask current = tasks.get(taskId);
if (current == null || current.status() != TaskStatus.PENDING
|| !current.assignee().equals(expectedAssignee)) {
return false;
}
tasks.put(taskId, new ApprovalTask(current.id(), current.instanceId(), current.stepId(), current.name(),
newAssignee, current.status(), current.createdAt(), current.completedAt(), current.action()));
return true;
}
@Override
public synchronized Page<ApprovalTask> query(TaskQuery query, PageRequest pageRequest) {
List<ApprovalTask> matched = tasks.values().stream().filter(task -> matches(task, query)).toList();
@@ -93,6 +93,64 @@ class InMemoryOrdoEngineTest {
assertThrows(TaskAlreadyCompletedException.class, () -> engine.approve(task.id(), "maria"));
}
@Test
void reassignsAPendingTaskToAnotherAssignee() {
var instance = engine.start("leave", "alice");
ApprovalTask task = engine.findPendingTasksByInstanceId(instance.id()).get(0);
ApprovalTask reassigned = engine.reassign(task.id(), "maria", "diana");
assertEquals(task.id(), reassigned.id());
assertEquals("diana", reassigned.assignee());
assertEquals(TaskStatus.PENDING, reassigned.status());
assertTrue(engine.findPendingTasksByAssignee("maria").isEmpty());
assertEquals(List.of(reassigned), engine.findPendingTasksByAssignee("diana"));
ProcessEvent event = engine.queryHistory(instance.id(), new PageRequest(0, 20)).content().stream()
.filter(e -> e.type() == ProcessEventType.TASK_REASSIGNED)
.findFirst()
.orElseThrow();
assertEquals("maria", event.actor());
assertEquals("diana", event.detail());
engine.approve(task.id(), "diana");
assertEquals("hr", engine.findPendingTasksByInstanceId(instance.id()).get(0).stepId());
}
@Test
void reassignIsRejectedForUnauthorizedCompletedSelfAndDuplicateAssignees() {
var instance = engine.start("leave", "alice");
ApprovalTask task = engine.findPendingTasksByInstanceId(instance.id()).get(0);
assertThrows(UnauthorizedTaskOperationException.class, () -> engine.reassign(task.id(), "mallory", "diana"));
assertThrows(IllegalArgumentException.class, () -> engine.reassign(task.id(), "maria", "maria"));
assertThrows(IllegalArgumentException.class, () -> engine.reassign(task.id(), "maria", " "));
engine.approve(task.id(), "maria");
assertThrows(TaskAlreadyCompletedException.class, () -> engine.reassign(task.id(), "maria", "diana"));
}
@Test
void reassignDoesNotAdvanceAnAllPolicyStepAndRejectsDuplicatePendingAssignee() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(ProcessDefinition.linear("leave-all-reassign", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL)
)));
var instance = allEngine.start("leave-all-reassign", "alice");
List<ApprovalTask> managerTasks = allEngine.findTasks(instance.id());
ApprovalTask mariaTask = managerTasks.stream().filter(t -> t.assignee().equals("maria")).findFirst().orElseThrow();
ApprovalTask mikeTask = managerTasks.stream().filter(t -> t.assignee().equals("mike")).findFirst().orElseThrow();
assertThrows(IllegalArgumentException.class, () -> allEngine.reassign(mariaTask.id(), "maria", "mike"));
ApprovalTask reassigned = allEngine.reassign(mariaTask.id(), "maria", "diana");
assertEquals(TaskStatus.PENDING, allEngine.findTask(mikeTask.id()).orElseThrow().status());
assertEquals(ProcessStatus.RUNNING, allEngine.findInstance(instance.id()).orElseThrow().status());
assertEquals(2, allEngine.findPendingTasksByInstanceId(instance.id()).size());
allEngine.approve(reassigned.id(), "diana");
assertEquals(ProcessStatus.RUNNING, allEngine.findInstance(instance.id()).orElseThrow().status());
allEngine.approve(mikeTask.id(), "mike");
assertEquals(ProcessStatus.APPROVED, allEngine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void exposesSpecificExceptionsForMissingAndDuplicateResources() {
assertThrows(DefinitionNotFoundException.class, () -> engine.start("missing", "alice"));
@@ -185,6 +243,7 @@ class InMemoryOrdoEngineTest {
ApprovalTask task = engine.findTasks(instance.id()).get(0);
assertThrows(IllegalArgumentException.class, () -> engine.approve(task.id(), " "));
assertThrows(IllegalArgumentException.class, () -> engine.reject(task.id(), " "));
assertThrows(IllegalArgumentException.class, () -> engine.reassign(task.id(), " ", "diana"));
}
@Test
@@ -75,6 +75,24 @@ class InMemoryApprovalTaskRepositoryTest {
assertEquals(List.of(leaveTask), byDefinition.content());
}
@Test
void reassignIfPendingUpdatesOnlyAMatchingPendingTask() {
InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository();
ApprovalTask pending = task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT);
repository.save(pending);
assertTrue(repository.reassignIfPending("task-1", "maria", "diana"));
assertEquals("diana", repository.findById("task-1").orElseThrow().assignee());
assertFalse(repository.reassignIfPending("task-1", "maria", "henry"));
assertEquals("diana", repository.findById("task-1").orElseThrow().assignee());
ApprovalTask completed = new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(), pending.name(),
"diana", TaskStatus.APPROVED, pending.createdAt(), CREATED_AT.plusSeconds(1), null);
assertTrue(repository.completeIfPending(completed));
assertFalse(repository.reassignIfPending("task-1", "diana", "henry"));
assertFalse(repository.reassignIfPending("missing", "maria", "diana"));
}
private static ApprovalTask task(String id, String instanceId, String assignee, TaskStatus status, Instant createdAt) {
return new ApprovalTask(id, instanceId, "step", "Step", assignee, status, createdAt, null, null);
}
@@ -27,7 +27,9 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
+ " VALUES (?, ?, ?, ?, ?, ?, ?)";
private static final String COMPLETE_IF_PENDING =
"UPDATE ordo_approval_task SET status = ?, completed_at = ?, action_actor = ?, action_comment = ?, action_at = ?"
+ " WHERE id = ? AND status = 'PENDING'";
+ " WHERE id = ? AND status = 'PENDING' AND assignee = ?";
private static final String REASSIGN_IF_PENDING =
"UPDATE ordo_approval_task SET assignee = ? WHERE id = ? AND status = 'PENDING' AND assignee = ?";
private static final String SELECT_TASK =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE id = ?";
private static final String SELECT_BY_INSTANCE =
@@ -119,6 +121,24 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
}
}
@Override
public boolean reassignIfPending(String taskId, String expectedAssignee, String newAssignee) {
Objects.requireNonNull(taskId, "taskId must not be null");
Objects.requireNonNull(expectedAssignee, "expectedAssignee must not be null");
Objects.requireNonNull(newAssignee, "newAssignee must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(REASSIGN_IF_PENDING)) {
update.setString(1, newAssignee);
update.setString(2, taskId);
update.setString(3, expectedAssignee);
return update.executeUpdate() == 1;
} catch (SQLException e) {
throw new JdbcStorageException("failed to reassign task: " + taskId, e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public Page<ApprovalTask> query(TaskQuery query, PageRequest pageRequest) {
Objects.requireNonNull(query, "query must not be null");
@@ -40,6 +40,7 @@ public final class ApprovalTaskMapper {
statement.setTimestamp(5, Timestamp.from(action.operatedAt()));
}
statement.setString(6, completedTask.id());
statement.setString(7, completedTask.assignee());
}
public static ApprovalTask read(ResultSet resultSet) throws SQLException {
@@ -154,6 +154,70 @@ class JdbcApprovalTaskRepositoryTest {
assertTrue(Set.of("first", "second").contains(stored.action().comment()));
}
@Test
void reassignIfPendingUpdatesAssigneeAndFailsWhenNotPending() {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
repository.save(pending);
assertTrue(repository.reassignIfPending("task-1", "maria", "diana"));
assertEquals("diana", repository.findById("task-1").orElseThrow().assignee());
assertEquals(List.of(repository.findById("task-1").orElseThrow()), repository.findPendingByAssignee("diana"));
assertTrue(repository.findPendingByAssignee("maria").isEmpty());
assertFalse(repository.reassignIfPending("task-1", "maria", "henry"));
assertTrue(repository.completeIfPending(completedTask(repository.findById("task-1").orElseThrow(), "ok")));
assertFalse(repository.reassignIfPending("task-1", "diana", "henry"));
}
@Test
void onlyOneOfConcurrentReassignAndCompletionWins() throws Exception {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
repository.save(pending);
CountDownLatch start = new CountDownLatch(1);
CountDownLatch done = new CountDownLatch(2);
AtomicInteger completeWins = new AtomicInteger();
AtomicInteger reassignWins = new AtomicInteger();
Thread completer = new Thread(() -> {
try {
start.await();
if (repository.completeIfPending(completedTask(pending, "ok"))) {
completeWins.incrementAndGet();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
});
Thread reassigner = new Thread(() -> {
try {
start.await();
if (repository.reassignIfPending("task-1", "maria", "diana")) {
reassignWins.incrementAndGet();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
});
completer.start();
reassigner.start();
start.countDown();
assertTrue(done.await(10, TimeUnit.SECONDS));
assertEquals(1, completeWins.get() + reassignWins.get());
ApprovalTask stored = repository.findById("task-1").orElseThrow();
if (completeWins.get() == 1) {
assertEquals(TaskStatus.APPROVED, stored.status());
assertEquals("maria", stored.assignee());
} else {
assertEquals(TaskStatus.PENDING, stored.status());
assertEquals("diana", stored.assignee());
}
}
private void attempt(CountDownLatch start, CountDownLatch done, AtomicInteger wins, ApprovalTask completed) {
try {
start.await();
@@ -97,6 +97,27 @@ class JdbcOrdoEngineIntegrationTest {
assertEquals("hr", engine.findPendingTasksByInstanceId(instance.id()).get(0).stepId());
}
@Test
void reassignsAPendingTaskAndRecordsHistory() {
ProcessInstance instance = engine.start("leave", "alice");
ApprovalTask task = engine.findPendingTasksByInstanceId(instance.id()).get(0);
ApprovalTask reassigned = engine.reassign(task.id(), "maria", "diana");
assertEquals("diana", reassigned.assignee());
assertEquals(TaskStatus.PENDING, reassigned.status());
assertTrue(engine.findPendingTasksByAssignee("maria").isEmpty());
assertEquals("diana", engine.findPendingTasksByAssignee("diana").get(0).assignee());
List<ProcessEvent> history = engine.queryHistory(instance.id(), new PageRequest(0, 20)).content();
ProcessEvent reassignedEvent = history.stream()
.filter(event -> event.type() == ProcessEventType.TASK_REASSIGNED)
.findFirst()
.orElseThrow();
assertEquals("maria", reassignedEvent.actor());
assertEquals("diana", reassignedEvent.detail());
assertEquals(task.id(), reassignedEvent.taskId());
}
@Test
void rollsBackTheWholeApprovalWhenTheNextStepCannotBeCreated() {
OrdoEngine failingEngine = newEngine((candidate, step, context) -> {