feat: support multi-approver steps with ANY/ALL policy

- ApprovalStep now has candidates + ApprovalPolicy (ANY/ALL) instead of a single assignee; ApprovalStep.single(...) kept as a convenience factory for the single-approver case
- AssigneeResolver resolves per-candidate; TaskStatus gains SKIPPED for auto-skipped sibling tasks
- DefaultOrdoEngine creates one task per candidate and advances per-step according to policy (ANY: first approval wins and skips the rest, only rejects once everyone rejects; ALL: fail-fast on first rejection, only advances once everyone approves)
- Added ApprovalTaskRepository.findByInstanceIdAndStepId and JDBC/in-memory implementations
- Added V2 migration: ordo_approval_step.policy column + ordo_step_candidate table
- Updated ApprovalStepMapper/ApprovalTaskMapper/JdbcProcessDefinitionRepository for the new schema, handling null action for SKIPPED tasks
- Migrated all call sites (tests, ordo-example, rhizome demo) to ApprovalStep.single(...) and the 3-arg AssigneeResolver
- Added ANY/ALL policy test coverage; fixed JdbcTestSupport to apply the V2 migration and strip SQL comments before splitting on ';'
This commit is contained in:
0264408
2026-09-10 18:20:33 +08:00
parent c524c74289
commit 581823290a
22 changed files with 440 additions and 90 deletions
@@ -1,5 +1,6 @@
package com.jetlumen.ordo.core;
import com.jetlumen.ordo.api.ApprovalPolicy;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.AssigneeResolver;
@@ -27,6 +28,7 @@ import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
/**
* Repository-backed implementation of the v0.1 linear approval runtime.
* Multistep write operations run inside a {@link TransactionExecutor} and
@@ -74,7 +76,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
ProcessInstance instance = new ProcessInstance(nextId(), definition.id(), initiator,
ProcessStatus.RUNNING, now, null, context);
instanceRepository.insert(instance);
createTask(instance, definition.steps().getFirst(), now);
createStepTasks(instance, definition.steps().getFirst(), now);
return instance;
});
}
@@ -85,14 +87,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
ApprovalTask task = requirePendingTaskForActor(taskId, actor);
Instant now = clock.instant();
ApprovalTask completedTask = completeTask(task, TaskStatus.APPROVED, new TaskAction(actor, comment, now));
ProcessInstance instance = requireInstance(task.instanceId());
ProcessDefinition definition = requireDefinition(instance.definitionId());
int stepIndex = indexOf(definition, task.stepId());
if (stepIndex == definition.steps().size() - 1) {
completeInstance(instance, ProcessStatus.APPROVED, now);
} else {
createTask(instance, definition.steps().get(stepIndex + 1), now);
}
advanceAfterDecision(completedTask, now);
return completedTask;
});
}
@@ -103,7 +98,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
ApprovalTask task = requirePendingTaskForActor(taskId, actor);
Instant now = clock.instant();
ApprovalTask completedTask = completeTask(task, TaskStatus.REJECTED, new TaskAction(actor, comment, now));
completeInstance(requireInstance(task.instanceId()), ProcessStatus.REJECTED, now);
advanceAfterDecision(completedTask, now);
return completedTask;
});
}
@@ -135,11 +130,13 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return taskRepository.findPendingByInstanceId(instanceId);
}
private void createTask(ProcessInstance instance, ApprovalStep step, Instant now) {
String assignee = assigneeResolver.resolve(step, instance.context());
requireText(assignee, "resolved assignee");
taskRepository.save(new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee,
TaskStatus.PENDING, now, null, null));
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now) {
for (String candidate : step.candidates()) {
String assignee = assigneeResolver.resolve(candidate, step, instance.context());
requireText(assignee, "resolved assignee");
taskRepository.save(new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee,
TaskStatus.PENDING, now, null, null));
}
}
private ApprovalTask requirePendingTaskForActor(String taskId, String actor) {
@@ -164,6 +161,68 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return completed;
}
/**
* Decides whether the step (and the process instance) can move on after a single candidate
* task was approved or rejected, applying the step's {@link ApprovalPolicy}.
*/
private void advanceAfterDecision(ApprovalTask completedTask, Instant now) {
ProcessInstance instance = requireInstance(completedTask.instanceId());
ProcessDefinition definition = requireDefinition(instance.definitionId());
ApprovalStep step = requireStep(definition, completedTask.stepId());
List<ApprovalTask> siblings = taskRepository.findByInstanceIdAndStepId(instance.id(), step.id());
if (completedTask.status() == TaskStatus.APPROVED && step.policy() == ApprovalPolicy.ANY) {
skipPendingSiblings(siblings, completedTask.id(), now);
advanceOrComplete(instance, definition, step, now);
return;
}
if (completedTask.status() == TaskStatus.REJECTED && step.policy() == ApprovalPolicy.ALL) {
skipPendingSiblings(siblings, completedTask.id(), now);
completeInstance(instance, ProcessStatus.REJECTED, now);
return;
}
if (completedTask.status() == TaskStatus.REJECTED) {
// ANY policy: the step is only rejected once every candidate has rejected it.
boolean stepStillAlive = siblings.stream()
.anyMatch(sibling -> !sibling.id().equals(completedTask.id())
&& (sibling.status() == TaskStatus.PENDING || sibling.status() == TaskStatus.APPROVED));
if (!stepStillAlive) {
completeInstance(instance, ProcessStatus.REJECTED, now);
}
return;
}
// ALL policy: only advance once every candidate has approved.
boolean allApproved = siblings.stream().allMatch(sibling -> sibling.status() == TaskStatus.APPROVED);
if (allApproved) {
advanceOrComplete(instance, definition, step, now);
}
}
private void advanceOrComplete(ProcessInstance instance, ProcessDefinition definition, ApprovalStep step, Instant now) {
int stepIndex = indexOf(definition, step.id());
if (stepIndex == definition.steps().size() - 1) {
completeInstance(instance, ProcessStatus.APPROVED, now);
} else {
createStepTasks(instance, definition.steps().get(stepIndex + 1), now);
}
}
/** Marks any still-pending sibling candidate tasks for the same step as skipped. */
private void skipPendingSiblings(List<ApprovalTask> siblings, String decidedTaskId, Instant now) {
for (ApprovalTask sibling : siblings) {
if (sibling.id().equals(decidedTaskId) || sibling.status() != TaskStatus.PENDING) {
continue;
}
ApprovalTask skipped = new ApprovalTask(sibling.id(), sibling.instanceId(), sibling.stepId(),
sibling.name(), sibling.assignee(), TaskStatus.SKIPPED, sibling.createdAt(), now, null);
// Best-effort: if another concurrent decision already completed this sibling, leave it as-is.
taskRepository.completeIfPending(skipped);
}
}
private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now) {
instanceRepository.update(new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
status, instance.startedAt(), now, instance.context()));
@@ -179,6 +238,13 @@ public final class DefaultOrdoEngine implements OrdoEngine {
.orElseThrow(() -> new IllegalStateException("instance not found: " + instanceId));
}
private static ApprovalStep requireStep(ProcessDefinition definition, String stepId) {
return definition.steps().stream()
.filter(step -> step.id().equals(stepId))
.findFirst()
.orElseThrow(() -> new IllegalStateException("step not found in definition: " + stepId));
}
private static int indexOf(ProcessDefinition definition, String stepId) {
for (int index = 0; index < definition.steps().size(); index++) {
if (definition.steps().get(index).id().equals(stepId)) {
@@ -28,6 +28,14 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
return tasks.values().stream().filter(task -> task.instanceId().equals(instanceId)).toList();
}
@Override
public synchronized List<ApprovalTask> findByInstanceIdAndStepId(String instanceId, String stepId) {
return tasks.values().stream()
.filter(task -> task.instanceId().equals(instanceId))
.filter(task -> task.stepId().equals(stepId))
.toList();
}
@Override
public synchronized List<ApprovalTask> findPendingByAssignee(String assignee) {
return tasks.values().stream()
@@ -1,5 +1,6 @@
package com.jetlumen.ordo.core;
import com.jetlumen.ordo.api.ApprovalPolicy;
import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.ProcessContext;
@@ -30,8 +31,8 @@ class InMemoryOrdoEngineTest {
void setUp() {
engine = new InMemoryOrdoEngine();
engine.register(new ProcessDefinition("leave", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry")
ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")
)));
}
@@ -79,7 +80,7 @@ class InMemoryOrdoEngineTest {
assertThrows(DefinitionNotFoundException.class, () -> engine.start("missing", "alice"));
assertThrows(TaskNotFoundException.class, () -> engine.approve("missing", "maria"));
assertThrows(DefinitionAlreadyExistsException.class, () -> engine.register(new ProcessDefinition(
"leave", "Another leave request", List.of(new ApprovalStep("lead", "Lead approval", "lee")))));
"leave", "Another leave request", List.of(ApprovalStep.single("lead", "Lead approval", "lee")))));
}
@Test
@@ -155,13 +156,13 @@ class InMemoryOrdoEngineTest {
@Test
void resolvesAssigneesFromTheProcessContext() {
InMemoryOrdoEngine contextAwareEngine = new InMemoryOrdoEngine((step, context) -> context.value(step.id())
InMemoryOrdoEngine contextAwareEngine = new InMemoryOrdoEngine((candidate, step, context) -> context.value(step.id())
.filter(String.class::isInstance)
.map(String.class::cast)
.orElse(step.assignee()));
.orElse(candidate));
contextAwareEngine.register(new ProcessDefinition("leave", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", "maria"),
new ApprovalStep("hr", "HR approval", "henry")
ApprovalStep.single("manager", "Manager approval", "maria"),
ApprovalStep.single("hr", "HR approval", "henry")
)));
var instance = contextAwareEngine.start("leave", "alice", new ProcessContext(Map.of(
@@ -175,4 +176,97 @@ class InMemoryOrdoEngineTest {
assertEquals("helena", contextAwareEngine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee());
assertEquals("david", instance.context().value("manager").orElseThrow());
}
@Test
void anyPolicyAdvancesOnFirstApprovalAndSkipsTheOtherCandidates() {
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
anyEngine.register(new ProcessDefinition("leave-any", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY),
ApprovalStep.single("hr", "HR approval", "henry")
)));
var instance = anyEngine.start("leave-any", "alice");
List<ApprovalTask> managerTasks = anyEngine.findTasks(instance.id());
assertEquals(2, managerTasks.size());
ApprovalTask mariaTask = managerTasks.stream().filter(t -> t.assignee().equals("maria")).findFirst().orElseThrow();
ApprovalTask mikeTask = managerTasks.stream().filter(t -> t.assignee().equals("mike")).findFirst().orElseThrow();
anyEngine.approve(mariaTask.id(), "maria");
assertEquals(TaskStatus.APPROVED, anyEngine.findTasks(instance.id()).stream()
.filter(t -> t.id().equals(mariaTask.id())).findFirst().orElseThrow().status());
assertEquals(TaskStatus.SKIPPED, anyEngine.findTasks(instance.id()).stream()
.filter(t -> t.id().equals(mikeTask.id())).findFirst().orElseThrow().status());
assertThrows(TaskAlreadyCompletedException.class, () -> anyEngine.approve(mikeTask.id(), "mike"));
ApprovalTask hrTask = anyEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertEquals("hr", hrTask.stepId());
anyEngine.approve(hrTask.id(), "henry");
assertEquals(ProcessStatus.APPROVED, anyEngine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void anyPolicyOnlyRejectsTheStepOnceEveryCandidateHasRejected() {
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
anyEngine.register(new ProcessDefinition("leave-any-reject", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ANY)
)));
var instance = anyEngine.start("leave-any-reject", "alice");
List<ApprovalTask> managerTasks = anyEngine.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();
anyEngine.reject(mariaTask.id(), "maria");
assertEquals(ProcessStatus.RUNNING, anyEngine.findInstance(instance.id()).orElseThrow().status());
assertEquals(TaskStatus.PENDING, anyEngine.findTask(mikeTask.id()).orElseThrow().status());
anyEngine.reject(mikeTask.id(), "mike");
assertEquals(ProcessStatus.REJECTED, anyEngine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void allPolicyOnlyAdvancesOnceEveryCandidateHasApproved() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(new ProcessDefinition("leave-all", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL),
ApprovalStep.single("hr", "HR approval", "henry")
)));
var instance = allEngine.start("leave-all", "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();
allEngine.approve(mariaTask.id(), "maria");
assertTrue(allEngine.findPendingTasksByInstanceId(instance.id()).stream()
.anyMatch(t -> t.id().equals(mikeTask.id())));
assertEquals(ProcessStatus.RUNNING, allEngine.findInstance(instance.id()).orElseThrow().status());
allEngine.approve(mikeTask.id(), "mike");
ApprovalTask hrTask = allEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertEquals("hr", hrTask.stepId());
allEngine.approve(hrTask.id(), "henry");
assertEquals(ProcessStatus.APPROVED, allEngine.findInstance(instance.id()).orElseThrow().status());
}
@Test
void allPolicyFailsFastAndSkipsRemainingCandidatesOnASingleRejection() {
InMemoryOrdoEngine allEngine = new InMemoryOrdoEngine();
allEngine.register(new ProcessDefinition("leave-all-reject", "Leave request", List.of(
new ApprovalStep("manager", "Manager approval", List.of("maria", "mike"), ApprovalPolicy.ALL)
)));
var instance = allEngine.start("leave-all-reject", "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();
allEngine.reject(mariaTask.id(), "maria");
assertEquals(ProcessStatus.REJECTED, allEngine.findInstance(instance.id()).orElseThrow().status());
assertEquals(TaskStatus.SKIPPED, allEngine.findTask(mikeTask.id()).orElseThrow().status());
assertNull(allEngine.findTask(mikeTask.id()).orElseThrow().action());
}
}