feat: add ACTION steps that run after the approval transaction commits
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionHandler;
|
||||
import com.jetlumen.ordo.api.ApprovalPolicy;
|
||||
import com.jetlumen.ordo.api.ApprovalStep;
|
||||
import com.jetlumen.ordo.api.ApprovalTask;
|
||||
@@ -10,6 +11,7 @@ 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.StepKind;
|
||||
import com.jetlumen.ordo.api.StepTransition;
|
||||
import com.jetlumen.ordo.api.TaskAction;
|
||||
import com.jetlumen.ordo.api.TaskStatus;
|
||||
@@ -30,6 +32,9 @@ import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
||||
|
||||
import java.time.Clock;
|
||||
import java.time.Instant;
|
||||
import java.lang.System.Logger;
|
||||
import java.lang.System.Logger.Level;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Comparator;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
@@ -44,22 +49,27 @@ import java.util.UUID;
|
||||
* when several JVMs share the same storage.
|
||||
*/
|
||||
public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
private static final int MAX_CONSECUTIVE_ACTIONS = 32;
|
||||
private static final Logger LOG = System.getLogger("ordo");
|
||||
|
||||
private final Clock clock;
|
||||
private final AssigneeResolver assigneeResolver;
|
||||
private final RoutingCondition routingCondition;
|
||||
private final ActionHandler actionHandler;
|
||||
private final TransactionExecutor transactionExecutor;
|
||||
private final ProcessDefinitionRepository definitionRepository;
|
||||
private final ProcessInstanceRepository instanceRepository;
|
||||
private final ApprovalTaskRepository taskRepository;
|
||||
|
||||
public DefaultOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
|
||||
TransactionExecutor transactionExecutor,
|
||||
ActionHandler actionHandler, TransactionExecutor transactionExecutor,
|
||||
ProcessDefinitionRepository definitionRepository,
|
||||
ProcessInstanceRepository instanceRepository,
|
||||
ApprovalTaskRepository taskRepository) {
|
||||
this.clock = Objects.requireNonNull(clock, "clock must not be null");
|
||||
this.assigneeResolver = Objects.requireNonNull(assigneeResolver, "assigneeResolver must not be null");
|
||||
this.routingCondition = Objects.requireNonNull(routingCondition, "routingCondition must not be null");
|
||||
this.actionHandler = Objects.requireNonNull(actionHandler, "actionHandler must not be null");
|
||||
this.transactionExecutor = Objects.requireNonNull(transactionExecutor, "transactionExecutor must not be null");
|
||||
this.definitionRepository = Objects.requireNonNull(definitionRepository, "definitionRepository must not be null");
|
||||
this.instanceRepository = Objects.requireNonNull(instanceRepository, "instanceRepository must not be null");
|
||||
@@ -93,26 +103,32 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
public synchronized ProcessInstance start(String definitionId, String initiator, ProcessContext context) {
|
||||
requireText(initiator, "initiator");
|
||||
Objects.requireNonNull(context, "context must not be null");
|
||||
return transactionExecutor.execute(() -> {
|
||||
List<PendingAction> queued = new ArrayList<>();
|
||||
ProcessInstance instance = transactionExecutor.execute(() -> {
|
||||
ProcessDefinition definition = requireDefinition(definitionId);
|
||||
Instant now = clock.instant();
|
||||
ProcessInstance instance = new ProcessInstance(nextId(), definition.id(), initiator,
|
||||
ProcessInstance started = new ProcessInstance(nextId(), definition.id(), initiator,
|
||||
ProcessStatus.RUNNING, now, null, context);
|
||||
instanceRepository.insert(instance);
|
||||
createStepTasks(instance, definition.steps().getFirst(), now);
|
||||
return instance;
|
||||
instanceRepository.insert(started);
|
||||
enterStep(started, definition, definition.steps().getFirst(), now, queued);
|
||||
return started;
|
||||
});
|
||||
runQueuedActions(queued);
|
||||
return instanceRepository.findById(instance.id()).orElse(instance);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized ApprovalTask approve(String taskId, String actor, String comment) {
|
||||
return transactionExecutor.execute(() -> {
|
||||
List<PendingAction> queued = new ArrayList<>();
|
||||
ApprovalTask completed = transactionExecutor.execute(() -> {
|
||||
ApprovalTask task = requirePendingTaskForActor(taskId, actor);
|
||||
Instant now = clock.instant();
|
||||
ApprovalTask completedTask = completeTask(task, TaskStatus.APPROVED, new TaskAction(actor, comment, now));
|
||||
advanceAfterDecision(completedTask, now);
|
||||
advanceAfterDecision(completedTask, now, queued);
|
||||
return completedTask;
|
||||
});
|
||||
runQueuedActions(queued);
|
||||
return completed;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -121,7 +137,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));
|
||||
advanceAfterDecision(completedTask, now);
|
||||
advanceAfterDecision(completedTask, now, List.of());
|
||||
return completedTask;
|
||||
});
|
||||
}
|
||||
@@ -217,7 +233,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
* 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) {
|
||||
private void advanceAfterDecision(ApprovalTask completedTask, Instant now, List<PendingAction> queued) {
|
||||
ProcessInstance instance = requireInstance(completedTask.instanceId());
|
||||
ProcessDefinition definition = requireDefinition(instance.definitionId());
|
||||
ApprovalStep step = requireStep(definition, completedTask.stepId());
|
||||
@@ -225,7 +241,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
|
||||
if (completedTask.status() == TaskStatus.APPROVED && step.policy() == ApprovalPolicy.ANY) {
|
||||
skipPendingSiblings(siblings, completedTask.id(), now);
|
||||
advanceOrComplete(instance, definition, step, now);
|
||||
advanceOrComplete(instance, definition, step, now, queued);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -249,16 +265,46 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
// 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);
|
||||
advanceOrComplete(instance, definition, step, now, queued);
|
||||
}
|
||||
}
|
||||
|
||||
private void advanceOrComplete(ProcessInstance instance, ProcessDefinition definition, ApprovalStep step, Instant now) {
|
||||
private void advanceOrComplete(ProcessInstance instance, ProcessDefinition definition, ApprovalStep step,
|
||||
Instant now, List<PendingAction> queued) {
|
||||
StepTransition matched = resolveTransition(definition, step, instance);
|
||||
if (matched.toStepId() == null) {
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now);
|
||||
} else {
|
||||
createStepTasks(instance, requireStep(definition, matched.toStepId()), now);
|
||||
enterStep(instance, definition, requireStep(definition, matched.toStepId()), now, queued);
|
||||
}
|
||||
}
|
||||
|
||||
private void enterStep(ProcessInstance instance, ProcessDefinition definition, ApprovalStep start, Instant now,
|
||||
List<PendingAction> queued) {
|
||||
ApprovalStep current = start;
|
||||
for (int hops = 0; hops < MAX_CONSECUTIVE_ACTIONS; hops++) {
|
||||
if (current.kind() == StepKind.APPROVAL) {
|
||||
createStepTasks(instance, current, now);
|
||||
return;
|
||||
}
|
||||
queued.add(new PendingAction(current.actionKey(), instance.context()));
|
||||
StepTransition matched = resolveTransition(definition, current, instance);
|
||||
if (matched.toStepId() == null) {
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now);
|
||||
return;
|
||||
}
|
||||
current = requireStep(definition, matched.toStepId());
|
||||
}
|
||||
throw new IllegalStateException("too many consecutive action steps in instance: " + instance.id());
|
||||
}
|
||||
|
||||
private void runQueuedActions(List<PendingAction> queued) {
|
||||
for (PendingAction pending : queued) {
|
||||
try {
|
||||
actionHandler.execute(pending.actionKey(), pending.context());
|
||||
} catch (RuntimeException e) {
|
||||
LOG.log(Level.WARNING, "action failed: " + pending.actionKey(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -324,4 +370,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
throw new IllegalArgumentException(name + " must not be blank");
|
||||
}
|
||||
}
|
||||
|
||||
private record PendingAction(String actionKey, ProcessContext context) {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionHandler;
|
||||
import com.jetlumen.ordo.api.ApprovalTask;
|
||||
import com.jetlumen.ordo.api.AssigneeResolver;
|
||||
import com.jetlumen.ordo.api.OrdoEngine;
|
||||
@@ -20,31 +21,37 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
|
||||
private final DefaultOrdoEngine delegate;
|
||||
|
||||
public InMemoryOrdoEngine() {
|
||||
this(Clock.systemUTC(), AssigneeResolver.direct(), RoutingCondition.always());
|
||||
this(Clock.systemUTC(), AssigneeResolver.direct(), RoutingCondition.always(), ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock) {
|
||||
this(clock, AssigneeResolver.direct(), RoutingCondition.always());
|
||||
this(clock, AssigneeResolver.direct(), RoutingCondition.always(), ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(AssigneeResolver assigneeResolver) {
|
||||
this(Clock.systemUTC(), assigneeResolver, RoutingCondition.always());
|
||||
this(Clock.systemUTC(), assigneeResolver, RoutingCondition.always(), ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver) {
|
||||
this(clock, assigneeResolver, RoutingCondition.always());
|
||||
this(clock, assigneeResolver, RoutingCondition.always(), ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(RoutingCondition routingCondition) {
|
||||
this(Clock.systemUTC(), AssigneeResolver.direct(), routingCondition);
|
||||
this(Clock.systemUTC(), AssigneeResolver.direct(), routingCondition, ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(AssigneeResolver assigneeResolver, RoutingCondition routingCondition) {
|
||||
this(Clock.systemUTC(), assigneeResolver, routingCondition);
|
||||
this(Clock.systemUTC(), assigneeResolver, routingCondition, ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition) {
|
||||
this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, new NoopTransactionExecutor(),
|
||||
this(clock, assigneeResolver, routingCondition, ActionHandler.noop());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
|
||||
ActionHandler actionHandler) {
|
||||
this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, actionHandler,
|
||||
new NoopTransactionExecutor(),
|
||||
new InMemoryProcessDefinitionRepository(),
|
||||
new InMemoryProcessInstanceRepository(),
|
||||
new InMemoryApprovalTaskRepository());
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.AssigneeResolver;
|
||||
import com.jetlumen.ordo.api.ApprovalPolicy;
|
||||
import com.jetlumen.ordo.api.ApprovalStep;
|
||||
import com.jetlumen.ordo.api.ApprovalTask;
|
||||
@@ -23,6 +24,7 @@ import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.time.Clock;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
@@ -386,6 +388,65 @@ class InMemoryOrdoEngineTest {
|
||||
assertTrue(routingEngine.findPendingTasksByInstanceId(low.id()).isEmpty());
|
||||
}
|
||||
|
||||
@Test
|
||||
void runsActionStepsAfterApprovalThenCreatesTheNextApprovalTask() {
|
||||
List<String> executed = new java.util.ArrayList<>();
|
||||
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> executed.add(key));
|
||||
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of(
|
||||
ApprovalStep.single("manager", "Manager approval", "maria"),
|
||||
ApprovalStep.action("notify", "Notify HR", "leave-approved-mail"),
|
||||
ApprovalStep.single("hr", "HR approval", "henry")),
|
||||
List.of(
|
||||
StepTransition.always("manager", "notify"),
|
||||
StepTransition.always("notify", "hr"),
|
||||
StepTransition.end("hr"))));
|
||||
|
||||
var instance = actionEngine.start("leave", "alice");
|
||||
actionEngine.approve(actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
|
||||
|
||||
assertEquals(List.of("leave-approved-mail"), executed);
|
||||
ApprovalTask hrTask = actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
|
||||
assertEquals("hr", hrTask.stepId());
|
||||
assertEquals(ProcessStatus.RUNNING, actionEngine.findInstance(instance.id()).orElseThrow().status());
|
||||
}
|
||||
|
||||
@Test
|
||||
void continuesWhenAnActionHandlerThrows() {
|
||||
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> {
|
||||
throw new IllegalStateException("mail failed");
|
||||
});
|
||||
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of(
|
||||
ApprovalStep.single("manager", "Manager approval", "maria"),
|
||||
ApprovalStep.action("notify", "Notify HR", "leave-approved-mail")),
|
||||
List.of(
|
||||
StepTransition.always("manager", "notify"),
|
||||
StepTransition.end("notify"))));
|
||||
|
||||
var instance = actionEngine.start("leave", "alice");
|
||||
actionEngine.approve(actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
|
||||
assertEquals(ProcessStatus.APPROVED, actionEngine.findInstance(instance.id()).orElseThrow().status());
|
||||
}
|
||||
|
||||
@Test
|
||||
void startEntersAnActionStepThenStopsOnTheFollowingApproval() {
|
||||
List<String> executed = new java.util.ArrayList<>();
|
||||
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> executed.add(key));
|
||||
actionEngine.register(new ProcessDefinition("leave", "Leave request", List.of(
|
||||
ApprovalStep.action("notify", "Notify manager", "leave-submitted-mail"),
|
||||
ApprovalStep.single("manager", "Manager approval", "maria")),
|
||||
List.of(
|
||||
StepTransition.always("notify", "manager"),
|
||||
StepTransition.end("manager"))));
|
||||
|
||||
var instance = actionEngine.start("leave", "alice");
|
||||
assertEquals(List.of("leave-submitted-mail"), executed);
|
||||
assertEquals("manager", actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst().stepId());
|
||||
assertEquals(ProcessStatus.RUNNING, instance.status());
|
||||
}
|
||||
|
||||
@Test
|
||||
void throwsWhenNoTransitionMatches() {
|
||||
InMemoryOrdoEngine routingEngine = new InMemoryOrdoEngine((key, context) -> false);
|
||||
|
||||
Reference in New Issue
Block a user