feat: persist process history events, listeners, and ACTION executions
Give hosts an append-only timeline, post-commit OrdoEventListener hooks, and durable ACTION results without blocking the approval flow. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1,13 +1,18 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionExecution;
|
||||
import com.jetlumen.ordo.api.ActionExecutionStatus;
|
||||
import com.jetlumen.ordo.api.ActionHandler;
|
||||
import com.jetlumen.ordo.api.ApprovalPolicy;
|
||||
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.OrdoEventListener;
|
||||
import com.jetlumen.ordo.api.ProcessContext;
|
||||
import com.jetlumen.ordo.api.ProcessDefinition;
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
@@ -30,8 +35,10 @@ 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.repository.ActionExecutionRepository;
|
||||
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
|
||||
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
|
||||
import com.jetlumen.ordo.api.repository.ProcessHistoryRepository;
|
||||
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
||||
|
||||
import java.time.Clock;
|
||||
@@ -44,6 +51,7 @@ import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
|
||||
/**
|
||||
@@ -64,12 +72,20 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
private final ProcessDefinitionRepository definitionRepository;
|
||||
private final ProcessInstanceRepository instanceRepository;
|
||||
private final ApprovalTaskRepository taskRepository;
|
||||
private final ProcessHistoryRepository historyRepository;
|
||||
private final ActionExecutionRepository actionExecutionRepository;
|
||||
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();
|
||||
|
||||
public DefaultOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
|
||||
ActionHandler actionHandler, TransactionExecutor transactionExecutor,
|
||||
ProcessDefinitionRepository definitionRepository,
|
||||
ProcessInstanceRepository instanceRepository,
|
||||
ApprovalTaskRepository taskRepository) {
|
||||
ApprovalTaskRepository taskRepository,
|
||||
ProcessHistoryRepository historyRepository,
|
||||
ActionExecutionRepository actionExecutionRepository,
|
||||
List<OrdoEventListener> listeners) {
|
||||
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");
|
||||
@@ -78,6 +94,10 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
this.definitionRepository = Objects.requireNonNull(definitionRepository, "definitionRepository must not be null");
|
||||
this.instanceRepository = Objects.requireNonNull(instanceRepository, "instanceRepository must not be null");
|
||||
this.taskRepository = Objects.requireNonNull(taskRepository, "taskRepository must not be null");
|
||||
this.historyRepository = Objects.requireNonNull(historyRepository, "historyRepository must not be null");
|
||||
this.actionExecutionRepository = Objects.requireNonNull(actionExecutionRepository,
|
||||
"actionExecutionRepository must not be null");
|
||||
this.listeners = List.copyOf(Objects.requireNonNull(listeners, "listeners must not be null"));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -108,49 +128,58 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
requireText(initiator, "initiator");
|
||||
Objects.requireNonNull(context, "context must not be null");
|
||||
List<PendingAction> queued = new ArrayList<>();
|
||||
List<ProcessEvent> events = new ArrayList<>();
|
||||
ProcessInstance instance = transactionExecutor.execute(() -> {
|
||||
ProcessDefinition definition = requireDefinition(definitionId);
|
||||
Instant now = clock.instant();
|
||||
ProcessInstance started = new ProcessInstance(nextId(), definition.id(), initiator,
|
||||
ProcessStatus.RUNNING, now, null, context);
|
||||
instanceRepository.insert(started);
|
||||
enterStep(started, definition, definition.steps().getFirst(), now, queued);
|
||||
record(events, started.id(), null, null, ProcessEventType.INSTANCE_STARTED, initiator, null, now);
|
||||
enterStep(started, definition, definition.steps().getFirst(), now, queued, events);
|
||||
return started;
|
||||
});
|
||||
runQueuedActions(queued);
|
||||
finishCommittedWork(queued, events);
|
||||
return instanceRepository.findById(instance.id()).orElse(instance);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized ApprovalTask approve(String taskId, String actor, String comment) {
|
||||
List<PendingAction> queued = new ArrayList<>();
|
||||
List<ProcessEvent> events = 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, queued);
|
||||
ApprovalTask completedTask = completeTask(task, TaskStatus.APPROVED, new TaskAction(actor, comment, now),
|
||||
events);
|
||||
advanceAfterDecision(completedTask, now, queued, events);
|
||||
return completedTask;
|
||||
});
|
||||
runQueuedActions(queued);
|
||||
finishCommittedWork(queued, events);
|
||||
return completed;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized ApprovalTask reject(String taskId, String actor, String comment) {
|
||||
return transactionExecutor.execute(() -> {
|
||||
List<ProcessEvent> events = new ArrayList<>();
|
||||
ApprovalTask completed = transactionExecutor.execute(() -> {
|
||||
ApprovalTask task = requirePendingTaskForActor(taskId, actor);
|
||||
Instant now = clock.instant();
|
||||
ApprovalTask completedTask = completeTask(task, TaskStatus.REJECTED, new TaskAction(actor, comment, now));
|
||||
advanceAfterDecision(completedTask, now, List.of());
|
||||
ApprovalTask completedTask = completeTask(task, TaskStatus.REJECTED, new TaskAction(actor, comment, now),
|
||||
events);
|
||||
advanceAfterDecision(completedTask, now, List.of(), events);
|
||||
return completedTask;
|
||||
});
|
||||
dispatch(events);
|
||||
return completed;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized ProcessInstance withdraw(String instanceId, String actor, String comment) {
|
||||
requireText(instanceId, "instance id");
|
||||
requireText(actor, "actor");
|
||||
return transactionExecutor.execute(() -> {
|
||||
List<ProcessEvent> events = new ArrayList<>();
|
||||
ProcessInstance withdrawn = transactionExecutor.execute(() -> {
|
||||
ProcessInstance instance = instanceRepository.findById(instanceId)
|
||||
.orElseThrow(() -> new InstanceNotFoundException(instanceId));
|
||||
if (!instance.initiator().equals(actor)) {
|
||||
@@ -160,19 +189,25 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
throw new InstanceAlreadyCompletedException(instanceId);
|
||||
}
|
||||
Instant now = clock.instant();
|
||||
ProcessInstance withdrawn = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
|
||||
ProcessInstance completed = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
|
||||
ProcessStatus.WITHDRAWN, instance.startedAt(), now, instance.context());
|
||||
if (!instanceRepository.completeIfRunning(withdrawn)) {
|
||||
if (!instanceRepository.completeIfRunning(completed)) {
|
||||
throw new InstanceAlreadyCompletedException(instanceId);
|
||||
}
|
||||
record(events, instanceId, null, null, ProcessEventType.INSTANCE_WITHDRAWN, 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);
|
||||
taskRepository.completeIfPending(skipped);
|
||||
if (taskRepository.completeIfPending(skipped)) {
|
||||
record(events, instanceId, pending.id(), pending.stepId(), ProcessEventType.TASK_SKIPPED, actor,
|
||||
comment, now);
|
||||
}
|
||||
}
|
||||
return withdrawn;
|
||||
return completed;
|
||||
});
|
||||
dispatch(events);
|
||||
return withdrawn;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -222,12 +257,28 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
return definitionRepository.findAll(pageRequest);
|
||||
}
|
||||
|
||||
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now) {
|
||||
@Override
|
||||
public synchronized Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) {
|
||||
requireText(instanceId, "instance id");
|
||||
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
|
||||
return historyRepository.query(instanceId, pageRequest);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized Page<ActionExecution> queryActionExecutions(String instanceId, PageRequest pageRequest) {
|
||||
requireText(instanceId, "instance id");
|
||||
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
|
||||
return actionExecutionRepository.query(instanceId, pageRequest);
|
||||
}
|
||||
|
||||
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now, List<ProcessEvent> events) {
|
||||
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));
|
||||
ApprovalTask task = new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee,
|
||||
TaskStatus.PENDING, now, null, null);
|
||||
taskRepository.save(task);
|
||||
record(events, instance.id(), task.id(), step.id(), ProcessEventType.TASK_CREATED, assignee, null, now);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -244,12 +295,17 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
return task;
|
||||
}
|
||||
|
||||
private ApprovalTask completeTask(ApprovalTask task, TaskStatus status, TaskAction action) {
|
||||
private ApprovalTask completeTask(ApprovalTask task, TaskStatus status, TaskAction action,
|
||||
List<ProcessEvent> events) {
|
||||
ApprovalTask completed = new ApprovalTask(task.id(), task.instanceId(), task.stepId(), task.name(),
|
||||
task.assignee(), status, task.createdAt(), action.operatedAt(), action);
|
||||
if (!taskRepository.completeIfPending(completed)) {
|
||||
throw new TaskAlreadyCompletedException(task.id());
|
||||
}
|
||||
ProcessEventType type = status == TaskStatus.APPROVED
|
||||
? ProcessEventType.TASK_APPROVED : ProcessEventType.TASK_REJECTED;
|
||||
record(events, task.instanceId(), task.id(), task.stepId(), type, action.actor(), action.comment(),
|
||||
action.operatedAt());
|
||||
return completed;
|
||||
}
|
||||
|
||||
@@ -257,21 +313,22 @@ 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, List<PendingAction> queued) {
|
||||
private void advanceAfterDecision(ApprovalTask completedTask, Instant now, List<PendingAction> queued,
|
||||
List<ProcessEvent> events) {
|
||||
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, queued);
|
||||
skipPendingSiblings(siblings, completedTask.id(), completedTask.action().actor(), now, events);
|
||||
advanceOrComplete(instance, definition, step, now, queued, events);
|
||||
return;
|
||||
}
|
||||
|
||||
if (completedTask.status() == TaskStatus.REJECTED && step.policy() == ApprovalPolicy.ALL) {
|
||||
skipPendingSiblings(siblings, completedTask.id(), now);
|
||||
completeInstance(instance, ProcessStatus.REJECTED, now);
|
||||
skipPendingSiblings(siblings, completedTask.id(), completedTask.action().actor(), now, events);
|
||||
completeInstance(instance, ProcessStatus.REJECTED, now, events);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -281,7 +338,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
.anyMatch(sibling -> !sibling.id().equals(completedTask.id())
|
||||
&& (sibling.status() == TaskStatus.PENDING || sibling.status() == TaskStatus.APPROVED));
|
||||
if (!stepStillAlive) {
|
||||
completeInstance(instance, ProcessStatus.REJECTED, now);
|
||||
completeInstance(instance, ProcessStatus.REJECTED, now, events);
|
||||
}
|
||||
return;
|
||||
}
|
||||
@@ -289,32 +346,37 @@ 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, queued);
|
||||
advanceOrComplete(instance, definition, step, now, queued, events);
|
||||
}
|
||||
}
|
||||
|
||||
private void advanceOrComplete(ProcessInstance instance, ProcessDefinition definition, ApprovalStep step,
|
||||
Instant now, List<PendingAction> queued) {
|
||||
Instant now, List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
StepTransition matched = resolveTransition(definition, step, instance);
|
||||
if (matched.toStepId() == null) {
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now);
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now, events);
|
||||
} else {
|
||||
enterStep(instance, definition, requireStep(definition, matched.toStepId()), now, queued);
|
||||
enterStep(instance, definition, requireStep(definition, matched.toStepId()), now, queued, events);
|
||||
}
|
||||
}
|
||||
|
||||
private void enterStep(ProcessInstance instance, ProcessDefinition definition, ApprovalStep start, Instant now,
|
||||
List<PendingAction> queued) {
|
||||
List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
ApprovalStep current = start;
|
||||
for (int hops = 0; hops < MAX_CONSECUTIVE_ACTIONS; hops++) {
|
||||
if (current.kind() == StepKind.APPROVAL) {
|
||||
createStepTasks(instance, current, now);
|
||||
createStepTasks(instance, current, now, events);
|
||||
return;
|
||||
}
|
||||
queued.add(new PendingAction(current.actionKey(), instance.context()));
|
||||
String executionId = nextId();
|
||||
actionExecutionRepository.insert(new ActionExecution(executionId, instance.id(), current.id(),
|
||||
current.actionKey(), ActionExecutionStatus.PENDING, null,
|
||||
now.plusMillis(eventSequence.getAndIncrement()), null));
|
||||
queued.add(new PendingAction(executionId, current.actionKey(), instance.id(), current.id(),
|
||||
instance.context()));
|
||||
StepTransition matched = resolveTransition(definition, current, instance);
|
||||
if (matched.toStepId() == null) {
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now);
|
||||
completeInstance(instance, ProcessStatus.APPROVED, now, events);
|
||||
return;
|
||||
}
|
||||
current = requireStep(definition, matched.toStepId());
|
||||
@@ -322,16 +384,49 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
throw new IllegalStateException("too many consecutive action steps in instance: " + instance.id());
|
||||
}
|
||||
|
||||
private void runQueuedActions(List<PendingAction> queued) {
|
||||
private void finishCommittedWork(List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
runQueuedActions(queued, events);
|
||||
dispatch(events);
|
||||
}
|
||||
|
||||
private void runQueuedActions(List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
for (PendingAction pending : queued) {
|
||||
Instant now = clock.instant();
|
||||
try {
|
||||
actionHandler.execute(pending.actionKey(), pending.context());
|
||||
actionExecutionRepository.complete(pending.executionId(), ActionExecutionStatus.SUCCESS, null, now);
|
||||
record(events, pending.instanceId(), null, pending.stepId(), ProcessEventType.ACTION_SUCCEEDED, null,
|
||||
pending.actionKey(), now);
|
||||
} catch (RuntimeException e) {
|
||||
actionExecutionRepository.complete(pending.executionId(), ActionExecutionStatus.FAILED, e.getMessage(),
|
||||
now);
|
||||
record(events, pending.instanceId(), null, pending.stepId(), ProcessEventType.ACTION_FAILED, null,
|
||||
e.getMessage(), now);
|
||||
LOG.log(Level.WARNING, "action failed: " + pending.actionKey(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void dispatch(List<ProcessEvent> events) {
|
||||
for (ProcessEvent event : events) {
|
||||
for (OrdoEventListener listener : listeners) {
|
||||
try {
|
||||
listener.onEvent(event);
|
||||
} catch (RuntimeException e) {
|
||||
LOG.log(Level.WARNING, "listener failed for event: " + event.type(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void record(List<ProcessEvent> events, String instanceId, String taskId, String stepId,
|
||||
ProcessEventType type, String actor, String detail, Instant occurredAt) {
|
||||
ProcessEvent event = new ProcessEvent(nextId(), instanceId, taskId, stepId, type, actor, detail,
|
||||
occurredAt.plusMillis(eventSequence.getAndIncrement()));
|
||||
historyRepository.append(event);
|
||||
events.add(event);
|
||||
}
|
||||
|
||||
private StepTransition resolveTransition(ProcessDefinition definition, ApprovalStep step, ProcessInstance instance) {
|
||||
List<StepTransition> candidates = definition.transitions().stream()
|
||||
.filter(transition -> transition.fromStepId().equals(step.id()))
|
||||
@@ -348,7 +443,8 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
}
|
||||
|
||||
/** Marks any still-pending sibling candidate tasks for the same step as skipped. */
|
||||
private void skipPendingSiblings(List<ApprovalTask> siblings, String decidedTaskId, Instant now) {
|
||||
private void skipPendingSiblings(List<ApprovalTask> siblings, String decidedTaskId, String actor, Instant now,
|
||||
List<ProcessEvent> events) {
|
||||
for (ApprovalTask sibling : siblings) {
|
||||
if (sibling.id().equals(decidedTaskId) || sibling.status() != TaskStatus.PENDING) {
|
||||
continue;
|
||||
@@ -356,16 +452,23 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
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);
|
||||
if (taskRepository.completeIfPending(skipped)) {
|
||||
record(events, sibling.instanceId(), sibling.id(), sibling.stepId(), ProcessEventType.TASK_SKIPPED,
|
||||
actor, null, now);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now) {
|
||||
private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now,
|
||||
List<ProcessEvent> events) {
|
||||
ProcessInstance completed = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
|
||||
status, instance.startedAt(), now, instance.context());
|
||||
if (!instanceRepository.completeIfRunning(completed)) {
|
||||
throw new InstanceAlreadyCompletedException(instance.id());
|
||||
}
|
||||
ProcessEventType type = status == ProcessStatus.APPROVED
|
||||
? ProcessEventType.INSTANCE_APPROVED : ProcessEventType.INSTANCE_REJECTED;
|
||||
record(events, instance.id(), null, null, type, null, null, now);
|
||||
}
|
||||
|
||||
private ProcessDefinition requireDefinition(String definitionId) {
|
||||
@@ -395,6 +498,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
}
|
||||
}
|
||||
|
||||
private record PendingAction(String actionKey, ProcessContext context) {
|
||||
private record PendingAction(String executionId, String actionKey, String instanceId, String stepId,
|
||||
ProcessContext context) {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,19 +1,24 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionExecution;
|
||||
import com.jetlumen.ordo.api.ActionHandler;
|
||||
import com.jetlumen.ordo.api.ApprovalTask;
|
||||
import com.jetlumen.ordo.api.AssigneeResolver;
|
||||
import com.jetlumen.ordo.api.OrdoEngine;
|
||||
import com.jetlumen.ordo.api.OrdoEventListener;
|
||||
import com.jetlumen.ordo.api.ProcessContext;
|
||||
import com.jetlumen.ordo.api.ProcessDefinition;
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
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.core.repository.InMemoryActionExecutionRepository;
|
||||
import com.jetlumen.ordo.core.repository.InMemoryApprovalTaskRepository;
|
||||
import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository;
|
||||
import com.jetlumen.ordo.core.repository.InMemoryProcessHistoryRepository;
|
||||
import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository;
|
||||
|
||||
import java.time.Clock;
|
||||
@@ -54,12 +59,20 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
|
||||
ActionHandler actionHandler) {
|
||||
this(clock, assigneeResolver, routingCondition, actionHandler, List.of());
|
||||
}
|
||||
|
||||
public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
|
||||
ActionHandler actionHandler, List<OrdoEventListener> listeners) {
|
||||
InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository();
|
||||
this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, actionHandler,
|
||||
new NoopTransactionExecutor(),
|
||||
new InMemoryProcessDefinitionRepository(),
|
||||
instanceRepository,
|
||||
new InMemoryApprovalTaskRepository(instanceRepository));
|
||||
new InMemoryApprovalTaskRepository(instanceRepository),
|
||||
new InMemoryProcessHistoryRepository(),
|
||||
new InMemoryActionExecutionRepository(),
|
||||
listeners);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -131,4 +144,14 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
|
||||
public Page<ProcessDefinition> listDefinitions(PageRequest pageRequest) {
|
||||
return delegate.listDefinitions(pageRequest);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Page<ProcessEvent> queryHistory(String instanceId, PageRequest pageRequest) {
|
||||
return delegate.queryHistory(instanceId, pageRequest);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Page<ActionExecution> queryActionExecutions(String instanceId, PageRequest pageRequest) {
|
||||
return delegate.queryActionExecutions(instanceId, pageRequest);
|
||||
}
|
||||
}
|
||||
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
package com.jetlumen.ordo.core.repository;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionExecution;
|
||||
import com.jetlumen.ordo.api.ActionExecutionStatus;
|
||||
import com.jetlumen.ordo.api.query.Page;
|
||||
import com.jetlumen.ordo.api.query.PageRequest;
|
||||
import com.jetlumen.ordo.api.repository.ActionExecutionRepository;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.Comparator;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
|
||||
/** Development-only in-memory implementation of the ACTION execution port. */
|
||||
public final class InMemoryActionExecutionRepository implements ActionExecutionRepository {
|
||||
private final Map<String, ActionExecution> executions = new LinkedHashMap<>();
|
||||
|
||||
@Override
|
||||
public synchronized void insert(ActionExecution execution) {
|
||||
Objects.requireNonNull(execution, "execution must not be null");
|
||||
executions.put(execution.id(), execution);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized boolean complete(String executionId, ActionExecutionStatus status, String errorMessage,
|
||||
Instant finishedAt) {
|
||||
ActionExecution current = executions.get(executionId);
|
||||
if (current == null || current.status() != ActionExecutionStatus.PENDING) {
|
||||
return false;
|
||||
}
|
||||
executions.put(executionId, new ActionExecution(current.id(), current.instanceId(), current.stepId(),
|
||||
current.actionKey(), status, errorMessage, current.startedAt(), finishedAt));
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized Page<ActionExecution> query(String instanceId, PageRequest pageRequest) {
|
||||
Objects.requireNonNull(instanceId, "instanceId must not be null");
|
||||
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
|
||||
List<ActionExecution> matched = executions.values().stream()
|
||||
.filter(execution -> execution.instanceId().equals(instanceId))
|
||||
.sorted(Comparator.comparing(ActionExecution::startedAt).thenComparing(ActionExecution::id))
|
||||
.toList();
|
||||
List<ActionExecution> page = matched.stream()
|
||||
.skip((long) pageRequest.offset())
|
||||
.limit(pageRequest.size())
|
||||
.toList();
|
||||
return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size());
|
||||
}
|
||||
}
|
||||
+37
@@ -0,0 +1,37 @@
|
||||
package com.jetlumen.ordo.core.repository;
|
||||
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.query.Page;
|
||||
import com.jetlumen.ordo.api.query.PageRequest;
|
||||
import com.jetlumen.ordo.api.repository.ProcessHistoryRepository;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Comparator;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
|
||||
/** Development-only in-memory implementation of the process history port. */
|
||||
public final class InMemoryProcessHistoryRepository implements ProcessHistoryRepository {
|
||||
private final List<ProcessEvent> events = new ArrayList<>();
|
||||
|
||||
@Override
|
||||
public synchronized void append(ProcessEvent event) {
|
||||
Objects.requireNonNull(event, "event must not be null");
|
||||
events.add(event);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized Page<ProcessEvent> query(String instanceId, PageRequest pageRequest) {
|
||||
Objects.requireNonNull(instanceId, "instanceId must not be null");
|
||||
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
|
||||
List<ProcessEvent> matched = events.stream()
|
||||
.filter(event -> event.instanceId().equals(instanceId))
|
||||
.sorted(Comparator.comparing(ProcessEvent::occurredAt).thenComparing(ProcessEvent::id))
|
||||
.toList();
|
||||
List<ProcessEvent> page = matched.stream()
|
||||
.skip((long) pageRequest.offset())
|
||||
.limit(pageRequest.size())
|
||||
.toList();
|
||||
return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size());
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,15 @@
|
||||
package com.jetlumen.ordo.core;
|
||||
|
||||
import com.jetlumen.ordo.api.ActionExecution;
|
||||
import com.jetlumen.ordo.api.ActionExecutionStatus;
|
||||
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.ProcessContext;
|
||||
import com.jetlumen.ordo.api.ProcessDefinition;
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
@@ -512,4 +516,106 @@ class InMemoryOrdoEngineTest {
|
||||
assertEquals(2, all.totalElements());
|
||||
assertEquals(List.of("expense", "leave"), all.content().stream().map(ProcessDefinition::id).toList());
|
||||
}
|
||||
|
||||
@Test
|
||||
void recordsHistoryForStartApproveSkipAndComplete() {
|
||||
InMemoryOrdoEngine anyEngine = new InMemoryOrdoEngine();
|
||||
anyEngine.register(ProcessDefinition.linear("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");
|
||||
ApprovalTask mariaTask = anyEngine.findTasks(instance.id()).stream()
|
||||
.filter(task -> task.assignee().equals("maria")).findFirst().orElseThrow();
|
||||
anyEngine.approve(mariaTask.id(), "maria");
|
||||
anyEngine.approve(anyEngine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "henry");
|
||||
|
||||
List<ProcessEventType> types = anyEngine.queryHistory(instance.id(), new PageRequest(0, 50)).content()
|
||||
.stream().map(ProcessEvent::type).toList();
|
||||
assertEquals(List.of(
|
||||
ProcessEventType.INSTANCE_STARTED,
|
||||
ProcessEventType.TASK_CREATED,
|
||||
ProcessEventType.TASK_CREATED,
|
||||
ProcessEventType.TASK_APPROVED,
|
||||
ProcessEventType.TASK_SKIPPED,
|
||||
ProcessEventType.TASK_CREATED,
|
||||
ProcessEventType.TASK_APPROVED,
|
||||
ProcessEventType.INSTANCE_APPROVED), types);
|
||||
}
|
||||
|
||||
@Test
|
||||
void recordsHistoryForWithdraw() {
|
||||
var instance = engine.start("leave", "alice");
|
||||
engine.withdraw(instance.id(), "alice", "changed plans");
|
||||
|
||||
List<ProcessEventType> types = engine.queryHistory(instance.id(), new PageRequest(0, 20)).content()
|
||||
.stream().map(ProcessEvent::type).toList();
|
||||
assertEquals(List.of(
|
||||
ProcessEventType.INSTANCE_STARTED,
|
||||
ProcessEventType.TASK_CREATED,
|
||||
ProcessEventType.INSTANCE_WITHDRAWN,
|
||||
ProcessEventType.TASK_SKIPPED), types);
|
||||
assertEquals("alice", engine.queryHistory(instance.id(), new PageRequest(0, 20)).content().get(2).actor());
|
||||
}
|
||||
|
||||
@Test
|
||||
void notifiesListenersAfterCommitAndIsolatesListenerFailures() {
|
||||
List<ProcessEventType> received = new java.util.ArrayList<>();
|
||||
InMemoryOrdoEngine listening = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> {
|
||||
}, List.of(event -> {
|
||||
if (event.type() == ProcessEventType.TASK_APPROVED) {
|
||||
throw new IllegalStateException("listener boom");
|
||||
}
|
||||
received.add(event.type());
|
||||
}));
|
||||
listening.register(ProcessDefinition.linear("leave", "Leave request", List.of(
|
||||
ApprovalStep.single("manager", "Manager approval", "maria"))));
|
||||
|
||||
var instance = listening.start("leave", "alice");
|
||||
listening.approve(listening.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
|
||||
|
||||
assertEquals(ProcessStatus.APPROVED, listening.findInstance(instance.id()).orElseThrow().status());
|
||||
assertEquals(List.of(
|
||||
ProcessEventType.INSTANCE_STARTED,
|
||||
ProcessEventType.TASK_CREATED,
|
||||
ProcessEventType.INSTANCE_APPROVED), received);
|
||||
}
|
||||
|
||||
@Test
|
||||
void persistsSuccessfulAndFailedActionExecutionsWithoutBlockingTheFlow() {
|
||||
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> {
|
||||
if (key.equals("fail-mail")) {
|
||||
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", "ok-mail"),
|
||||
ApprovalStep.action("fail", "Fail mail", "fail-mail")),
|
||||
List.of(
|
||||
StepTransition.always("manager", "notify"),
|
||||
StepTransition.always("notify", "fail"),
|
||||
StepTransition.end("fail"))));
|
||||
|
||||
var instance = actionEngine.start("leave", "alice");
|
||||
actionEngine.approve(actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
|
||||
|
||||
assertEquals(ProcessStatus.APPROVED, actionEngine.findInstance(instance.id()).orElseThrow().status());
|
||||
List<ActionExecution> executions = actionEngine.queryActionExecutions(instance.id(), new PageRequest(0, 10))
|
||||
.content();
|
||||
assertEquals(2, executions.size());
|
||||
assertEquals("ok-mail", executions.get(0).actionKey());
|
||||
assertEquals(ActionExecutionStatus.SUCCESS, executions.get(0).status());
|
||||
assertEquals("fail-mail", executions.get(1).actionKey());
|
||||
assertEquals(ActionExecutionStatus.FAILED, executions.get(1).status());
|
||||
assertEquals("mail failed", executions.get(1).errorMessage());
|
||||
List<ProcessEventType> types = actionEngine.queryHistory(instance.id(), new PageRequest(0, 50)).content()
|
||||
.stream().map(ProcessEvent::type).toList();
|
||||
assertTrue(types.contains(ProcessEventType.ACTION_SUCCEEDED));
|
||||
assertTrue(types.contains(ProcessEventType.ACTION_FAILED));
|
||||
assertTrue(types.contains(ProcessEventType.INSTANCE_APPROVED));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user