From f3220fd0a002cc0be962cb30eb105154ad5a6b0e Mon Sep 17 00:00:00 2001 From: 0264408 Date: Tue, 15 Sep 2026 10:14:44 +0800 Subject: [PATCH] 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 --- docs/roadmap.md | 51 ++--- .../jetlumen/ordo/api/ActionExecution.java | 25 +++ .../ordo/api/ActionExecutionStatus.java | 8 + .../com/jetlumen/ordo/api/OrdoEngine.java | 6 + .../jetlumen/ordo/api/OrdoEventListener.java | 10 + .../com/jetlumen/ordo/api/ProcessEvent.java | 29 +++ .../jetlumen/ordo/api/ProcessEventType.java | 15 ++ .../repository/ActionExecutionRepository.java | 23 +++ .../repository/ProcessHistoryRepository.java | 13 ++ .../jetlumen/ordo/core/DefaultOrdoEngine.java | 178 ++++++++++++++---- .../ordo/core/InMemoryOrdoEngine.java | 25 ++- .../InMemoryActionExecutionRepository.java | 52 +++++ .../InMemoryProcessHistoryRepository.java | 37 ++++ .../ordo/core/InMemoryOrdoEngineTest.java | 106 +++++++++++ .../spring/OrdoJdbcAutoConfiguration.java | 26 ++- .../spring/OrdoJdbcAutoConfigurationTest.java | 42 +++++ .../jdbc/JdbcActionExecutionRepository.java | 94 +++++++++ .../jdbc/JdbcProcessHistoryRepository.java | 73 +++++++ .../jdbc/mapper/ActionExecutionMapper.java | 50 +++++ .../jdbc/mapper/ProcessEventMapper.java | 47 +++++ ...add_process_event_and_action_execution.sql | 29 +++ .../JdbcActionExecutionRepositoryTest.java | 66 +++++++ .../jdbc/JdbcOrdoEngineIntegrationTest.java | 57 +++++- .../jdbc/JdbcPostgresIntegrationTest.java | 5 +- .../JdbcProcessHistoryRepositoryTest.java | 60 ++++++ .../ordo/storage/jdbc/JdbcTestSupport.java | 3 +- 26 files changed, 1063 insertions(+), 67 deletions(-) create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecution.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecutionStatus.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEventListener.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEvent.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ActionExecutionRepository.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessHistoryRepository.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryActionExecutionRepository.java create mode 100644 ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessHistoryRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ActionExecutionMapper.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessEventMapper.java create mode 100644 ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepositoryTest.java create mode 100644 ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepositoryTest.java diff --git a/docs/roadmap.md b/docs/roadmap.md index 3a969fa..893e9fe 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -3,38 +3,43 @@ 记录当前已完成能力之后,后续要做的开发计划。按优先级分组,供后续排期/立项参考。 ## 已完成(背景,非本文档重点) + - ANY/ALL 多候选人会签/或签(见 `/memories/repo/any-all-multi-approval.md`) - JDBC 存储模块 + Flyway 迁移(V1~V4) - 条件路由 `StepTransition` + `RoutingCondition` - ACTION 步骤 + `ActionHandler` - 流程撤回(`WITHDRAWN`) - **任务/实例/流程定义分页过滤查询 API**(2026-09-15 完成):`ApprovalTaskRepository.query`、 - `ProcessInstanceRepository.query`、`ProcessDefinitionRepository.findAll` + `OrdoEngine` 对应的 - `queryTasks`/`queryInstances`/`listDefinitions`,均支持 `PageRequest`/`Page` 分页与按 - assignee/instanceId/definitionId/status/initiator/时间范围过滤,默认按时间降序(最新优先)。 +`ProcessInstanceRepository.query`、`ProcessDefinitionRepository.findAll` + `OrdoEngine` 对应的 +`queryTasks`/`queryInstances`/`listDefinitions`,均支持 `PageRequest`/`Page` 分页与按 +assignee/instanceId/definitionId/status/initiator/时间范围过滤,默认按时间降序(最新优先)。 ## P0 — 审计与扩展点 -1. **独立历史/审计事件模型**:现在历史只能靠 `ApprovalTask.action` 字段拼凑,没有独立的流程事件表 - (谁在何时对哪个实例做了什么)。建议新增 `ProcessEvent`/`ProcessHistoryRepository`。 -2. **状态变更事件监听器**:`OrdoEngine` 目前没有任何 listener/hook,无法在任务创建、审批、实例完成时 - 被外部感知(做通知、写审计日志等)。可加 `OrdoEventListener` 扩展点,风格与 `AssigneeResolver`/ - `RoutingCondition` 一致。 -3. **ACTION 步骤执行记录持久化**:目前 `ActionHandler` 执行结果只在宿主内存里记(如 rhizome 的 - `LeaveActionHandler`),重启即丢失,且失败只打日志不影响流程状态,需要设计重试/失败处理策略。 + +- [x] **独立历史/审计事件模型**:现在历史只能靠 `ApprovalTask.action` 字段拼凑,没有独立的流程事件表 + (谁在何时对哪个实例做了什么)。建议新增 `ProcessEvent`/`ProcessHistoryRepository`。 +- [x] **状态变更事件监听器**:`OrdoEngine` 目前没有任何 listener/hook,无法在任务创建、审批、实例完成时 + 被外部感知(做通知、写审计日志等)。可加 `OrdoEventListener` 扩展点,风格与 `AssigneeResolver`/ + `RoutingCondition` 一致。 +- [x] **ACTION 步骤执行记录持久化**:目前 `ActionHandler` 执行结果只在宿主内存里记(如 rhizome 的 + `LeaveActionHandler`),重启即丢失,且失败只打日志不影响流程状态,需要设计重试/失败处理策略。 ## P1 — 任务生命周期完善 -4. **任务委托/转派(delegate/reassign)**:任务创建后 assignee 不可变,无法转交他人处理。 -5. **超时/升级(SLA/escalation)**:无到期时间、定时器、自动升级机制。 -6. **流程实例取消 vs 撤回**:目前只有 `WITHDRAWN`(仅发起人可操作),没有管理员/系统层面的 - `CANCELLED` 语义。 + +- [ ] **任务委托/转派(delegate/reassign)**:任务创建后 assignee 不可变,无法转交他人处理。 +- [ ] **超时/升级(SLA/escalation)**:无到期时间、定时器、自动升级机制。 +- [ ] **流程实例取消 vs 撤回**:目前只有 `WITHDRAWN`(仅发起人可操作),没有管理员/系统层面的 + `CANCELLED` 语义。 ## P2 — 架构级演进(范围较大,放在后面) -7. **流程定义版本化**:目前同 id 直接整体替换(`replace`),建议演进为不可变多版本 + 运行中实例 - 锁定所用版本。 -8. **多租户支持**:数据模型无 tenant 隔离字段。 -9. **JDBC 多方言支持**:目前 DDL/实现明显偏向 PostgreSQL(唯一键冲突处理等),无 MySQL/Testcontainers - 测试,若要支持更多数据库需要抽象 dialect 层。 -10. **通用 REST Starter**:现在 REST 层完全是 rhizome 自己写的 demo,可考虑提供一个可选的 - `ordo-spring-boot-starter-web` 暴露标准 REST 接口。 -11. **表单/UI schema、子流程、并行 fork-join**:属于更大的引擎能力扩展,优先级最低,等基础能力稳定后 - 再评估是否需要。 + +- [ ] **流程定义版本化**:目前同 id 直接整体替换(`replace`),建议演进为不可变多版本 + 运行中实例 + 锁定所用版本。 +- [ ] **多租户支持**:数据模型无 tenant 隔离字段。 +- [ ] **JDBC 多方言支持**:目前 DDL/实现明显偏向 PostgreSQL(唯一键冲突处理等),无 MySQL/Testcontainers + 测试,若要支持更多数据库需要抽象 dialect 层。 +- [ ] **通用 REST Starter**:现在 REST 层完全是 rhizome 自己写的 demo,可考虑提供一个可选的 + `ordo-spring-boot-starter-web` 暴露标准 REST 接口。 +- [ ] **表单/UI schema、子流程、并行 fork-join**:属于更大的引擎能力扩展,优先级最低,等基础能力稳定后 + 再评估是否需要。 + diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecution.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecution.java new file mode 100644 index 0000000..745df5a --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecution.java @@ -0,0 +1,25 @@ +package com.jetlumen.ordo.api; + +import java.time.Instant; +import java.util.Objects; + +/** Persisted record of an ACTION step invocation. */ +public record ActionExecution(String id, String instanceId, String stepId, String actionKey, + ActionExecutionStatus status, String errorMessage, Instant startedAt, + Instant finishedAt) { + public ActionExecution { + requireText(id, "id"); + requireText(instanceId, "instanceId"); + requireText(stepId, "stepId"); + requireText(actionKey, "actionKey"); + Objects.requireNonNull(status, "status must not be null"); + Objects.requireNonNull(startedAt, "startedAt must not be null"); + errorMessage = errorMessage == null || errorMessage.isBlank() ? null : errorMessage.strip(); + } + + private static void requireText(String value, String name) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(name + " must not be blank"); + } + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecutionStatus.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecutionStatus.java new file mode 100644 index 0000000..08bfc00 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ActionExecutionStatus.java @@ -0,0 +1,8 @@ +package com.jetlumen.ordo.api; + +/** Lifecycle of a persisted ACTION-step execution. */ +public enum ActionExecutionStatus { + PENDING, + SUCCESS, + FAILED +} 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 ae5736e..5971a2f 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 @@ -47,4 +47,10 @@ public interface OrdoEngine { /** Paginated listing of all registered process definitions. */ Page listDefinitions(PageRequest pageRequest); + + /** Instance timeline, oldest-first; see {@link com.jetlumen.ordo.api.repository.ProcessHistoryRepository#query}. */ + Page queryHistory(String instanceId, PageRequest pageRequest); + + /** ACTION executions for an instance, oldest-first. */ + Page queryActionExecutions(String instanceId, PageRequest pageRequest); } diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEventListener.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEventListener.java new file mode 100644 index 0000000..6dd1f13 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/OrdoEventListener.java @@ -0,0 +1,10 @@ +package com.jetlumen.ordo.api; + +/** + * Host hook invoked after a process mutation has been committed. Implementations must not throw + * in a way that affects the engine: the runtime isolates listener failures. + */ +@FunctionalInterface +public interface OrdoEventListener { + void onEvent(ProcessEvent event); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEvent.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEvent.java new file mode 100644 index 0000000..9ac5f4e --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEvent.java @@ -0,0 +1,29 @@ +package com.jetlumen.ordo.api; + +import java.time.Instant; +import java.util.Objects; + +/** Immutable audit record of something that happened to a process instance. */ +public record ProcessEvent(String id, String instanceId, String taskId, String stepId, ProcessEventType type, + String actor, String detail, Instant occurredAt) { + public ProcessEvent { + requireText(id, "id"); + requireText(instanceId, "instanceId"); + Objects.requireNonNull(type, "type must not be null"); + Objects.requireNonNull(occurredAt, "occurredAt must not be null"); + taskId = blankToNull(taskId); + stepId = blankToNull(stepId); + actor = blankToNull(actor); + detail = blankToNull(detail); + } + + private static String blankToNull(String value) { + return value == null || value.isBlank() ? null : value.strip(); + } + + private static void requireText(String value, String name) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(name + " must not be blank"); + } + } +} 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 new file mode 100644 index 0000000..b952e15 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/ProcessEventType.java @@ -0,0 +1,15 @@ +package com.jetlumen.ordo.api; + +/** Kinds of append-only process history events. */ +public enum ProcessEventType { + INSTANCE_STARTED, + TASK_CREATED, + TASK_APPROVED, + TASK_REJECTED, + TASK_SKIPPED, + INSTANCE_APPROVED, + INSTANCE_REJECTED, + INSTANCE_WITHDRAWN, + ACTION_SUCCEEDED, + ACTION_FAILED +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ActionExecutionRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ActionExecutionRepository.java new file mode 100644 index 0000000..ca92498 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ActionExecutionRepository.java @@ -0,0 +1,23 @@ +package com.jetlumen.ordo.api.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 java.time.Instant; + +/** Storage port for ACTION-step execution records. */ +public interface ActionExecutionRepository { + void insert(ActionExecution execution); + + /** + * Completes a pending execution. + * + * @return true if the row was still pending and was updated + */ + boolean complete(String executionId, ActionExecutionStatus status, String errorMessage, Instant finishedAt); + + /** Executions for one instance, oldest-first ({@code started_at}, then {@code id}). */ + Page query(String instanceId, PageRequest pageRequest); +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessHistoryRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessHistoryRepository.java new file mode 100644 index 0000000..a59ff08 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessHistoryRepository.java @@ -0,0 +1,13 @@ +package com.jetlumen.ordo.api.repository; + +import com.jetlumen.ordo.api.ProcessEvent; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; + +/** Storage port for append-only process history. */ +public interface ProcessHistoryRepository { + void append(ProcessEvent event); + + /** Timeline for one instance, oldest-first ({@code occurred_at}, then {@code id}). */ + Page query(String instanceId, 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 07fa326..6827841 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 @@ -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 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 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 queued = new ArrayList<>(); + List 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 queued = new ArrayList<>(); + List 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 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 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 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 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 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 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 queued) { + private void advanceAfterDecision(ApprovalTask completedTask, Instant now, List queued, + List events) { ProcessInstance instance = requireInstance(completedTask.instanceId()); ProcessDefinition definition = requireDefinition(instance.definitionId()); ApprovalStep step = requireStep(definition, completedTask.stepId()); List 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 queued) { + Instant now, List queued, List 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 queued) { + List queued, List 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 queued) { + private void finishCommittedWork(List queued, List events) { + runQueuedActions(queued, events); + dispatch(events); + } + + private void runQueuedActions(List queued, List 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 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 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 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 siblings, String decidedTaskId, Instant now) { + private void skipPendingSiblings(List siblings, String decidedTaskId, String actor, Instant now, + List 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 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) { } } 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 937e74f..ac20939 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 @@ -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 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 listDefinitions(PageRequest pageRequest) { return delegate.listDefinitions(pageRequest); } + + @Override + public Page queryHistory(String instanceId, PageRequest pageRequest) { + return delegate.queryHistory(instanceId, pageRequest); + } + + @Override + public Page queryActionExecutions(String instanceId, PageRequest pageRequest) { + return delegate.queryActionExecutions(instanceId, pageRequest); + } } diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryActionExecutionRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryActionExecutionRepository.java new file mode 100644 index 0000000..1f93b8b --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryActionExecutionRepository.java @@ -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 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 query(String instanceId, PageRequest pageRequest) { + Objects.requireNonNull(instanceId, "instanceId must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + List matched = executions.values().stream() + .filter(execution -> execution.instanceId().equals(instanceId)) + .sorted(Comparator.comparing(ActionExecution::startedAt).thenComparing(ActionExecution::id)) + .toList(); + List page = matched.stream() + .skip((long) pageRequest.offset()) + .limit(pageRequest.size()) + .toList(); + return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size()); + } +} diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessHistoryRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessHistoryRepository.java new file mode 100644 index 0000000..619fc69 --- /dev/null +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessHistoryRepository.java @@ -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 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 query(String instanceId, PageRequest pageRequest) { + Objects.requireNonNull(instanceId, "instanceId must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + List matched = events.stream() + .filter(event -> event.instanceId().equals(instanceId)) + .sorted(Comparator.comparing(ProcessEvent::occurredAt).thenComparing(ProcessEvent::id)) + .toList(); + List page = matched.stream() + .skip((long) pageRequest.offset()) + .limit(pageRequest.size()) + .toList(); + return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size()); + } +} 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 9597cae..127b281 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 @@ -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 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 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 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 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 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)); + } } diff --git a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java index ded1321..9a2d0ca 100644 --- a/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java +++ b/ordo-spring-boot-autoconfigure/src/main/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfiguration.java @@ -3,17 +3,23 @@ package com.jetlumen.ordo.spring; import com.jetlumen.ordo.api.ActionHandler; import com.jetlumen.ordo.api.AssigneeResolver; import com.jetlumen.ordo.api.OrdoEngine; +import com.jetlumen.ordo.api.OrdoEventListener; import com.jetlumen.ordo.api.RoutingCondition; import com.jetlumen.ordo.api.TransactionExecutor; +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 com.jetlumen.ordo.core.DefaultOrdoEngine; +import com.jetlumen.ordo.storage.jdbc.JdbcActionExecutionRepository; import com.jetlumen.ordo.storage.jdbc.JdbcApprovalTaskRepository; import com.jetlumen.ordo.storage.jdbc.JdbcConnectionProvider; import com.jetlumen.ordo.storage.jdbc.JdbcProcessDefinitionRepository; +import com.jetlumen.ordo.storage.jdbc.JdbcProcessHistoryRepository; import com.jetlumen.ordo.storage.jdbc.JdbcProcessInstanceRepository; import com.jetlumen.ordo.storage.jdbc.JdbcTransactionExecutor; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; @@ -103,6 +109,18 @@ public class OrdoJdbcAutoConfiguration { return new JdbcApprovalTaskRepository(connectionProvider); } + @Bean + @ConditionalOnMissingBean + public ProcessHistoryRepository ordoProcessHistoryRepository(JdbcConnectionProvider connectionProvider) { + return new JdbcProcessHistoryRepository(connectionProvider); + } + + @Bean + @ConditionalOnMissingBean + public ActionExecutionRepository ordoActionExecutionRepository(JdbcConnectionProvider connectionProvider) { + return new JdbcActionExecutionRepository(connectionProvider); + } + @Bean @ConditionalOnMissingBean public OrdoEngine ordoEngine(Clock ordoClock, @@ -112,10 +130,14 @@ public class OrdoJdbcAutoConfiguration { TransactionExecutor ordoTransactionExecutor, ProcessDefinitionRepository ordoProcessDefinitionRepository, ProcessInstanceRepository ordoProcessInstanceRepository, - ApprovalTaskRepository ordoApprovalTaskRepository) { + ApprovalTaskRepository ordoApprovalTaskRepository, + ProcessHistoryRepository ordoProcessHistoryRepository, + ActionExecutionRepository ordoActionExecutionRepository, + ObjectProvider ordoEventListeners) { return new DefaultOrdoEngine(ordoClock, ordoAssigneeResolver, ordoRoutingCondition, actionHandler, ordoTransactionExecutor, ordoProcessDefinitionRepository, ordoProcessInstanceRepository, - ordoApprovalTaskRepository); + ordoApprovalTaskRepository, ordoProcessHistoryRepository, ordoActionExecutionRepository, + ordoEventListeners.orderedStream().toList()); } @Bean diff --git a/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java b/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java index d553ba0..e982f58 100644 --- a/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java +++ b/ordo-spring-boot-autoconfigure/src/test/java/com/jetlumen/ordo/spring/OrdoJdbcAutoConfigurationTest.java @@ -5,9 +5,13 @@ 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.ProcessDefinition; +import com.jetlumen.ordo.api.ProcessEvent; +import com.jetlumen.ordo.api.ProcessEventType; import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.RoutingCondition; +import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; @@ -151,6 +155,27 @@ class OrdoJdbcAutoConfigurationTest { }); } + @Test + void invokesOrdoEventListenerBeans() { + withDataSourceRunner.withUserConfiguration(RecordingListenerConfig.class) + .run(context -> { + OrdoEngine engine = context.getBean(OrdoEngine.class); + engine.register(LEAVE_REQUEST); + ProcessInstance instance = engine.start("leave-request", "alice"); + engine.approve(engine.findPendingTasksByAssignee("maria").getFirst().id(), "maria"); + + RecordingListener listener = context.getBean(RecordingListener.class); + assertThat(listener.types).contains( + ProcessEventType.INSTANCE_STARTED, + ProcessEventType.TASK_CREATED, + ProcessEventType.TASK_APPROVED, + ProcessEventType.INSTANCE_APPROVED); + assertThat(listener.types.getFirst()).isEqualTo(ProcessEventType.INSTANCE_STARTED); + assertThat(engine.queryHistory(instance.id(), new PageRequest(0, 20)) + .totalElements()).isGreaterThan(0); + }); + } + @Configuration static class CustomAssigneeResolverConfig { @Bean @@ -176,4 +201,21 @@ class OrdoJdbcAutoConfigurationTest { }; } } + + static class RecordingListener implements OrdoEventListener { + final List types = new java.util.concurrent.CopyOnWriteArrayList<>(); + + @Override + public void onEvent(ProcessEvent event) { + types.add(event.type()); + } + } + + @Configuration + static class RecordingListenerConfig { + @Bean + RecordingListener recordingListener() { + return new RecordingListener(); + } + } } diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java new file mode 100644 index 0000000..0e95127 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepository.java @@ -0,0 +1,94 @@ +package com.jetlumen.ordo.storage.jdbc; + +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 com.jetlumen.ordo.storage.jdbc.mapper.ActionExecutionMapper; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; + +/** JDBC implementation of the ACTION execution port. */ +public final class JdbcActionExecutionRepository implements ActionExecutionRepository { + private static final String COLUMNS = + "id, instance_id, step_id, action_key, status, error_message, started_at, finished_at"; + private static final String INSERT = + "INSERT INTO ordo_action_execution (id, instance_id, step_id, action_key, status, started_at)" + + " VALUES (?, ?, ?, ?, ?, ?)"; + private static final String COMPLETE_IF_PENDING = + "UPDATE ordo_action_execution SET status = ?, error_message = ?, finished_at = ?" + + " WHERE id = ? AND status = 'PENDING'"; + + private final JdbcConnectionProvider connectionProvider; + + public JdbcActionExecutionRepository(JdbcConnectionProvider connectionProvider) { + this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + } + + @Override + public void insert(ActionExecution execution) { + Objects.requireNonNull(execution, "execution must not be null"); + Connection connection = connectionProvider.getConnection(); + try (PreparedStatement insert = connection.prepareStatement(INSERT)) { + ActionExecutionMapper.bindInsert(insert, execution); + insert.executeUpdate(); + } catch (SQLException e) { + throw new JdbcStorageException("failed to insert action execution: " + execution.id(), e); + } finally { + connectionProvider.close(connection); + } + } + + @Override + public boolean complete(String executionId, ActionExecutionStatus status, String errorMessage, Instant finishedAt) { + Objects.requireNonNull(executionId, "executionId must not be null"); + Objects.requireNonNull(status, "status must not be null"); + Objects.requireNonNull(finishedAt, "finishedAt must not be null"); + Connection connection = connectionProvider.getConnection(); + try (PreparedStatement update = connection.prepareStatement(COMPLETE_IF_PENDING)) { + ActionExecutionMapper.bindComplete(update, executionId, status, errorMessage, finishedAt); + return update.executeUpdate() == 1; + } catch (SQLException e) { + throw new JdbcStorageException("failed to complete action execution: " + executionId, e); + } finally { + connectionProvider.close(connection); + } + } + + @Override + public Page query(String instanceId, PageRequest pageRequest) { + Objects.requireNonNull(instanceId, "instanceId must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + Connection connection = connectionProvider.getConnection(); + try { + long total = PageSupport.count(connection, + "SELECT COUNT(*) FROM ordo_action_execution WHERE instance_id = ?", List.of(instanceId)); + String sql = "SELECT " + COLUMNS + " FROM ordo_action_execution WHERE instance_id = ?" + + " ORDER BY started_at ASC, id ASC LIMIT ? OFFSET ?"; + List content = new ArrayList<>(); + try (PreparedStatement select = connection.prepareStatement(sql)) { + select.setString(1, instanceId); + select.setInt(2, pageRequest.size()); + select.setInt(3, pageRequest.offset()); + try (ResultSet resultSet = select.executeQuery()) { + while (resultSet.next()) { + content.add(ActionExecutionMapper.read(resultSet)); + } + } + } + return new Page<>(content, total, pageRequest.page(), pageRequest.size()); + } catch (SQLException e) { + throw new JdbcStorageException("failed to query action executions for instance: " + instanceId, e); + } finally { + connectionProvider.close(connection); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java new file mode 100644 index 0000000..b135eee --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepository.java @@ -0,0 +1,73 @@ +package com.jetlumen.ordo.storage.jdbc; + +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 com.jetlumen.ordo.storage.jdbc.mapper.ProcessEventMapper; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; + +/** JDBC implementation of the process history port. */ +public final class JdbcProcessHistoryRepository implements ProcessHistoryRepository { + private static final String COLUMNS = + "id, instance_id, task_id, step_id, event_type, actor, detail, occurred_at"; + private static final String INSERT = + "INSERT INTO ordo_process_event (id, instance_id, task_id, step_id, event_type, actor, detail, occurred_at)" + + " VALUES (?, ?, ?, ?, ?, ?, ?, ?)"; + + private final JdbcConnectionProvider connectionProvider; + + public JdbcProcessHistoryRepository(JdbcConnectionProvider connectionProvider) { + this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); + } + + @Override + public void append(ProcessEvent event) { + Objects.requireNonNull(event, "event must not be null"); + Connection connection = connectionProvider.getConnection(); + try (PreparedStatement insert = connection.prepareStatement(INSERT)) { + ProcessEventMapper.bindInsert(insert, event); + insert.executeUpdate(); + } catch (SQLException e) { + throw new JdbcStorageException("failed to append process event: " + event.id(), e); + } finally { + connectionProvider.close(connection); + } + } + + @Override + public Page query(String instanceId, PageRequest pageRequest) { + Objects.requireNonNull(instanceId, "instanceId must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + Connection connection = connectionProvider.getConnection(); + try { + long total = PageSupport.count(connection, + "SELECT COUNT(*) FROM ordo_process_event WHERE instance_id = ?", List.of(instanceId)); + String sql = "SELECT " + COLUMNS + " FROM ordo_process_event WHERE instance_id = ?" + + " ORDER BY occurred_at ASC, id ASC LIMIT ? OFFSET ?"; + List content = new ArrayList<>(); + try (PreparedStatement select = connection.prepareStatement(sql)) { + select.setString(1, instanceId); + select.setInt(2, pageRequest.size()); + select.setInt(3, pageRequest.offset()); + try (ResultSet resultSet = select.executeQuery()) { + while (resultSet.next()) { + content.add(ProcessEventMapper.read(resultSet)); + } + } + } + return new Page<>(content, total, pageRequest.page(), pageRequest.size()); + } catch (SQLException e) { + throw new JdbcStorageException("failed to query process events for instance: " + instanceId, e); + } finally { + connectionProvider.close(connection); + } + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ActionExecutionMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ActionExecutionMapper.java new file mode 100644 index 0000000..5dc2dbe --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ActionExecutionMapper.java @@ -0,0 +1,50 @@ +package com.jetlumen.ordo.storage.jdbc.mapper; + +import com.jetlumen.ordo.api.ActionExecution; +import com.jetlumen.ordo.api.ActionExecutionStatus; + +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Timestamp; +import java.sql.Types; + +/** Maps rows of {@code ordo_action_execution} to {@link ActionExecution} objects and back. */ +public final class ActionExecutionMapper { + private ActionExecutionMapper() { + } + + public static void bindInsert(PreparedStatement statement, ActionExecution execution) throws SQLException { + statement.setString(1, execution.id()); + statement.setString(2, execution.instanceId()); + statement.setString(3, execution.stepId()); + statement.setString(4, execution.actionKey()); + statement.setString(5, execution.status().name()); + statement.setTimestamp(6, Timestamp.from(execution.startedAt())); + } + + public static ActionExecution read(ResultSet resultSet) throws SQLException { + Timestamp finishedAt = resultSet.getTimestamp("finished_at"); + return new ActionExecution( + resultSet.getString("id"), + resultSet.getString("instance_id"), + resultSet.getString("step_id"), + resultSet.getString("action_key"), + ActionExecutionStatus.valueOf(resultSet.getString("status")), + resultSet.getString("error_message"), + resultSet.getTimestamp("started_at").toInstant(), + finishedAt == null ? null : finishedAt.toInstant()); + } + + public static void bindComplete(PreparedStatement statement, String executionId, ActionExecutionStatus status, + String errorMessage, java.time.Instant finishedAt) throws SQLException { + statement.setString(1, status.name()); + if (errorMessage == null || errorMessage.isBlank()) { + statement.setNull(2, Types.VARCHAR); + } else { + statement.setString(2, errorMessage); + } + statement.setTimestamp(3, Timestamp.from(finishedAt)); + statement.setString(4, executionId); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessEventMapper.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessEventMapper.java new file mode 100644 index 0000000..ef330c0 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/mapper/ProcessEventMapper.java @@ -0,0 +1,47 @@ +package com.jetlumen.ordo.storage.jdbc.mapper; + +import com.jetlumen.ordo.api.ProcessEvent; +import com.jetlumen.ordo.api.ProcessEventType; + +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Timestamp; +import java.sql.Types; + +/** Maps rows of {@code ordo_process_event} to {@link ProcessEvent} objects and back. */ +public final class ProcessEventMapper { + private ProcessEventMapper() { + } + + public static void bindInsert(PreparedStatement statement, ProcessEvent event) throws SQLException { + statement.setString(1, event.id()); + statement.setString(2, event.instanceId()); + setNullableString(statement, 3, event.taskId()); + setNullableString(statement, 4, event.stepId()); + statement.setString(5, event.type().name()); + setNullableString(statement, 6, event.actor()); + setNullableString(statement, 7, event.detail()); + statement.setTimestamp(8, Timestamp.from(event.occurredAt())); + } + + public static ProcessEvent read(ResultSet resultSet) throws SQLException { + return new ProcessEvent( + resultSet.getString("id"), + resultSet.getString("instance_id"), + resultSet.getString("task_id"), + resultSet.getString("step_id"), + ProcessEventType.valueOf(resultSet.getString("event_type")), + resultSet.getString("actor"), + resultSet.getString("detail"), + resultSet.getTimestamp("occurred_at").toInstant()); + } + + private static void setNullableString(PreparedStatement statement, int index, String value) throws SQLException { + if (value == null) { + statement.setNull(index, Types.VARCHAR); + } else { + statement.setString(index, value); + } + } +} diff --git a/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql b/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql new file mode 100644 index 0000000..52ed75d --- /dev/null +++ b/ordo-storage-jdbc/src/main/resources/db/migration/V5__add_process_event_and_action_execution.sql @@ -0,0 +1,29 @@ +-- Process history events and ACTION-step execution records. + +CREATE TABLE ordo_process_event ( + id VARCHAR(36) PRIMARY KEY, + instance_id VARCHAR(36) NOT NULL, + task_id VARCHAR(36), + step_id VARCHAR(64), + event_type VARCHAR(32) NOT NULL, + actor VARCHAR(255), + detail TEXT, + occurred_at TIMESTAMP NOT NULL, + CONSTRAINT fk_process_event_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +); + +CREATE INDEX idx_process_event_instance_time ON ordo_process_event (instance_id, occurred_at, id); + +CREATE TABLE ordo_action_execution ( + id VARCHAR(36) PRIMARY KEY, + instance_id VARCHAR(36) NOT NULL, + step_id VARCHAR(64) NOT NULL, + action_key VARCHAR(255) NOT NULL, + status VARCHAR(32) NOT NULL, + error_message TEXT, + started_at TIMESTAMP NOT NULL, + finished_at TIMESTAMP, + CONSTRAINT fk_action_execution_instance FOREIGN KEY (instance_id) REFERENCES ordo_process_instance (id) +); + +CREATE INDEX idx_action_execution_instance ON ordo_action_execution (instance_id); diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepositoryTest.java new file mode 100644 index 0000000..6f6b5c3 --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcActionExecutionRepositoryTest.java @@ -0,0 +1,66 @@ +package com.jetlumen.ordo.storage.jdbc; + +import com.jetlumen.ordo.api.ActionExecution; +import com.jetlumen.ordo.api.ActionExecutionStatus; +import com.jetlumen.ordo.api.ApprovalStep; +import com.jetlumen.ordo.api.ProcessContext; +import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.ProcessStatus; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class JdbcActionExecutionRepositoryTest { + private static final Instant T0 = Instant.parse("2026-03-01T08:00:00Z"); + + private JdbcActionExecutionRepository repository; + + @BeforeEach + void setUp() { + JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); + repository = new JdbcActionExecutionRepository(connectionProvider); + new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", + "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); + new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", "alice", + ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); + } + + @Test + void roundTripsPendingThenCompletedExecution() { + ActionExecution pending = new ActionExecution("ex-1", "inst-1", "notify", "leave-mail", + ActionExecutionStatus.PENDING, null, T0, null); + repository.insert(pending); + + assertTrue(repository.complete("ex-1", ActionExecutionStatus.SUCCESS, null, T0.plusSeconds(2))); + Page page = repository.query("inst-1", new PageRequest(0, 10)); + assertEquals(1, page.totalElements()); + ActionExecution stored = page.content().getFirst(); + assertEquals(ActionExecutionStatus.SUCCESS, stored.status()); + assertEquals(T0.plusSeconds(2), stored.finishedAt()); + assertFalse(repository.complete("ex-1", ActionExecutionStatus.FAILED, "nope", T0.plusSeconds(3))); + } + + @Test + void completeFailsForUnknownExecution() { + assertFalse(repository.complete("missing", ActionExecutionStatus.FAILED, "x", T0)); + } + + @Test + void queryOrdersOldestFirst() { + repository.insert(new ActionExecution("ex-1", "inst-1", "a", "first", ActionExecutionStatus.PENDING, null, + T0, null)); + repository.insert(new ActionExecution("ex-2", "inst-1", "b", "second", ActionExecutionStatus.PENDING, null, + T0.plusSeconds(1), null)); + Page page = repository.query("inst-1", new PageRequest(0, 10)); + assertEquals(List.of("ex-1", "ex-2"), page.content().stream().map(ActionExecution::id).toList()); + } +} 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 271a1c0..a4fc069 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 @@ -1,5 +1,6 @@ package com.jetlumen.ordo.storage.jdbc; +import com.jetlumen.ordo.api.ActionExecutionStatus; import com.jetlumen.ordo.api.ActionHandler; import com.jetlumen.ordo.api.ApprovalStep; import com.jetlumen.ordo.api.ApprovalTask; @@ -7,8 +8,11 @@ import com.jetlumen.ordo.api.AssigneeResolver; import com.jetlumen.ordo.api.OrdoEngine; 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.query.PageRequest; import com.jetlumen.ordo.api.RoutingCondition; import com.jetlumen.ordo.api.StepTransition; import com.jetlumen.ordo.api.TaskStatus; @@ -124,7 +128,10 @@ class JdbcOrdoEngineIntegrationTest { new JdbcTransactionExecutor(connectionProvider), new JdbcProcessDefinitionRepository(connectionProvider), new JdbcProcessInstanceRepository(connectionProvider), - new JdbcApprovalTaskRepository(connectionProvider)); + new JdbcApprovalTaskRepository(connectionProvider), + new JdbcProcessHistoryRepository(connectionProvider), + new JdbcActionExecutionRepository(connectionProvider), + List.of()); failingEngine.register(new ProcessDefinition("leave-noroute", "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")), List.of(StepTransition.endWhen("manager", "never", 0)))); @@ -137,6 +144,9 @@ class JdbcOrdoEngineIntegrationTest { assertEquals(TaskStatus.PENDING, storedTask.status()); assertNull(storedTask.action()); assertEquals(ProcessStatus.RUNNING, failingEngine.findInstance(instance.id()).orElseThrow().status()); + assertEquals(List.of(ProcessEventType.INSTANCE_STARTED, ProcessEventType.TASK_CREATED), + failingEngine.queryHistory(instance.id(), new PageRequest(0, 20)).content().stream() + .map(ProcessEvent::type).toList()); } @Test @@ -196,12 +206,55 @@ class JdbcOrdoEngineIntegrationTest { assertThrows(InstanceAlreadyCompletedException.class, () -> engine.withdraw(instance.id(), "alice")); } + @Test + void persistHistoryAndActionExecutions() { + OrdoEngine actionEngine = new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), AssigneeResolver.direct(), + RoutingCondition.always(), + (key, context) -> { + if ("boom".equals(key)) { + throw new IllegalStateException("mail failed"); + } + }, + new JdbcTransactionExecutor(connectionProvider), + new JdbcProcessDefinitionRepository(connectionProvider), + new JdbcProcessInstanceRepository(connectionProvider), + new JdbcApprovalTaskRepository(connectionProvider), + new JdbcProcessHistoryRepository(connectionProvider), + new JdbcActionExecutionRepository(connectionProvider), + List.of()); + actionEngine.register(new ProcessDefinition("leave-action", "Leave request", List.of( + ApprovalStep.single("manager", "Manager approval", "maria"), + ApprovalStep.action("notify", "Notify", "ok-mail"), + ApprovalStep.action("fail", "Fail", "boom")), + List.of( + StepTransition.always("manager", "notify"), + StepTransition.always("notify", "fail"), + StepTransition.end("fail")))); + + ProcessInstance instance = actionEngine.start("leave-action", "alice"); + actionEngine.approve(actionEngine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria"); + + assertEquals(ProcessStatus.APPROVED, actionEngine.findInstance(instance.id()).orElseThrow().status()); + assertEquals(ActionExecutionStatus.SUCCESS, + actionEngine.queryActionExecutions(instance.id(), new PageRequest(0, 10)).content().getFirst().status()); + assertEquals(ActionExecutionStatus.FAILED, + actionEngine.queryActionExecutions(instance.id(), new PageRequest(0, 10)).content().get(1).status()); + List 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)); + } + private OrdoEngine newEngine(AssigneeResolver assigneeResolver) { return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver, RoutingCondition.always(), ActionHandler.noop(), new JdbcTransactionExecutor(connectionProvider), new JdbcProcessDefinitionRepository(connectionProvider), new JdbcProcessInstanceRepository(connectionProvider), - new JdbcApprovalTaskRepository(connectionProvider)); + new JdbcApprovalTaskRepository(connectionProvider), + new JdbcProcessHistoryRepository(connectionProvider), + new JdbcActionExecutionRepository(connectionProvider), + List.of()); } } diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java index 4a894a6..5ea1bd7 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcPostgresIntegrationTest.java @@ -224,7 +224,10 @@ class JdbcPostgresIntegrationTest { new JdbcTransactionExecutor(connectionProvider), new JdbcProcessDefinitionRepository(connectionProvider), new JdbcProcessInstanceRepository(connectionProvider), - new JdbcApprovalTaskRepository(connectionProvider)); + new JdbcApprovalTaskRepository(connectionProvider), + new JdbcProcessHistoryRepository(connectionProvider), + new JdbcActionExecutionRepository(connectionProvider), + List.of()); } private static PGSimpleDataSource newDataSource(String currentSchema) { diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepositoryTest.java new file mode 100644 index 0000000..618d938 --- /dev/null +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessHistoryRepositoryTest.java @@ -0,0 +1,60 @@ +package com.jetlumen.ordo.storage.jdbc; + +import com.jetlumen.ordo.api.ApprovalStep; +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.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class JdbcProcessHistoryRepositoryTest { + private static final Instant T0 = Instant.parse("2026-03-01T08:00:00Z"); + + private JdbcProcessHistoryRepository repository; + + @BeforeEach + void setUp() { + JdbcConnectionProvider connectionProvider = new JdbcConnectionProvider(JdbcTestSupport.newDataSource()); + repository = new JdbcProcessHistoryRepository(connectionProvider); + new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("leave", + "Leave request", List.of(ApprovalStep.single("manager", "Manager approval", "maria")))); + new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-1", "leave", "alice", + ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); + new JdbcProcessInstanceRepository(connectionProvider).insert(new ProcessInstance("inst-2", "leave", "bob", + ProcessStatus.RUNNING, T0, null, ProcessContext.empty())); + } + + @Test + void appendsAndQueriesOldestFirstForOneInstance() { + ProcessEvent started = event("e1", "inst-1", ProcessEventType.INSTANCE_STARTED, T0); + ProcessEvent created = event("e2", "inst-1", ProcessEventType.TASK_CREATED, T0.plusSeconds(1)); + ProcessEvent other = event("e3", "inst-2", ProcessEventType.INSTANCE_STARTED, T0); + repository.append(started); + repository.append(created); + repository.append(other); + + Page page = repository.query("inst-1", new PageRequest(0, 10)); + assertEquals(2, page.totalElements()); + assertEquals(List.of(started, created), page.content()); + + Page first = repository.query("inst-1", new PageRequest(0, 1)); + assertEquals(1, first.content().size()); + assertEquals("e1", first.content().getFirst().id()); + assertTrue(first.hasNext()); + } + + private static ProcessEvent event(String id, String instanceId, ProcessEventType type, Instant at) { + return new ProcessEvent(id, instanceId, null, null, type, "alice", null, at); + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java index 3f35c29..bbb8f06 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcTestSupport.java @@ -17,7 +17,8 @@ final class JdbcTestSupport { "/db/migration/V1__create_ordo_tables.sql", "/db/migration/V2__add_step_candidates_and_policy.sql", "/db/migration/V3__add_step_transitions.sql", - "/db/migration/V4__add_step_kind_and_action_key.sql" + "/db/migration/V4__add_step_kind_and_action_key.sql", + "/db/migration/V5__add_process_event_and_action_execution.sql" }; private static final String[] SCHEMA_SQL = loadSchemas();