From 850583f32625c08631bcec88fc38ff5c01199ef4 Mon Sep 17 00:00:00 2001 From: 0264408 Date: Wed, 16 Sep 2026 08:40:25 +0800 Subject: [PATCH] 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 --- README.md | 12 ++-- docs/roadmap.md | 23 +++++-- docs/usage.md | 27 ++++---- .../com/jetlumen/ordo/api/OrdoEngine.java | 1 + .../jetlumen/ordo/api/ProcessEventType.java | 1 + .../repository/ApprovalTaskRepository.java | 13 +++- .../jetlumen/ordo/core/DefaultOrdoEngine.java | 30 +++++++++ .../ordo/core/InMemoryOrdoEngine.java | 5 ++ .../InMemoryApprovalTaskRepository.java | 15 ++++- .../ordo/core/InMemoryOrdoEngineTest.java | 59 +++++++++++++++++ .../InMemoryApprovalTaskRepositoryTest.java | 18 ++++++ .../jdbc/JdbcApprovalTaskRepository.java | 22 ++++++- .../jdbc/mapper/ApprovalTaskMapper.java | 1 + .../jdbc/JdbcApprovalTaskRepositoryTest.java | 64 +++++++++++++++++++ .../jdbc/JdbcOrdoEngineIntegrationTest.java | 21 ++++++ 15 files changed, 286 insertions(+), 26 deletions(-) diff --git a/README.md b/README.md index c29f973..4c4a9f5 100644 --- a/README.md +++ b/README.md @@ -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。 ## 查询 diff --git a/docs/roadmap.md b/docs/roadmap.md index 56385db..a8002b0 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -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)。 diff --git a/docs/usage.md b/docs/usage.md index 84576e6..52d1da0 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -10,9 +10,9 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。 做:流程定义、实例推进、待办任务、条件路由、ACTION 副作用、审计事件、分页查询。 -产品边界(不做,且暂不在开发计划):业务表单、用户体系、官方 REST、多租户。业务字段放在 `ProcessContext`(不可变 `Map`)。 +产品边界:不做业务表单、用户体系、多租户;不内置设计器 UI。业务字段放在 `ProcessContext`(不可变 `Map`)。当前也**没有** 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 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 由宿主自建。 diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java index 76d3368..3309261 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEngine.java @@ -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); } diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java index b952e15..462c75f 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java @@ -7,6 +7,7 @@ public enum ProcessEventType { TASK_APPROVED, TASK_REJECTED, TASK_SKIPPED, + TASK_REASSIGNED, INSTANCE_APPROVED, INSTANCE_REJECTED, INSTANCE_WITHDRAWN, diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java index 7729f84..830eb03 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java @@ -25,12 +25,21 @@ public interface ApprovalTaskRepository { List 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 query(TaskQuery query, PageRequest pageRequest); } diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java index 60e478b..fe2cae7 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/DefaultOrdoEngine.java @@ -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 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"); diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java index f33819c..b69fc56 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/InMemoryOrdoEngine.java @@ -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); diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java index 2389011..e5529cf 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java @@ -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 query(TaskQuery query, PageRequest pageRequest) { List matched = tasks.values().stream().filter(task -> matches(task, query)).toList(); diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java index 6c44d1c..187fa83 100644 --- a/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/InMemoryOrdoEngineTest.java @@ -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 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 diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java index 107f23a..fdd654f 100644 --- a/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java @@ -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); } diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java index 0103f01..b39fd57 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java @@ -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 query(TaskQuery query, PageRequest pageRequest) { Objects.requireNonNull(query, "query must not be null"); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java index 86b279c..ce93d71 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ApprovalTaskMapper.java @@ -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 { diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java index a024bf4..513a737 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java @@ -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(); diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java index 3e0a9ca..eba2aac 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcOrdoEngineIntegrationTest.java @@ -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 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) -> {