From 77d1e5198b54c46f09b74dee6ac92432776cfe9a Mon Sep 17 00:00:00 2001 From: 0264408 Date: Tue, 15 Sep 2026 08:56:05 +0800 Subject: [PATCH] feat: Implement pagination and filtering for approval tasks and process instances - Added query method to InMemoryApprovalTaskRepository for filtering and pagination of approval tasks. - Enhanced InMemoryProcessInstanceRepository with query method for filtering and pagination of process instances. - Introduced Page and PageRequest classes for handling paginated results. - Updated JdbcApprovalTaskRepository and JdbcProcessInstanceRepository to support querying with pagination. - Added tests for querying, filtering, and pagination in InMemory and JDBC repositories. - Implemented dynamic WHERE clause building for SQL queries in PageSupport. --- docs/roadmap.md | 40 +++++++ .../com/jetlumen/ordo/api/OrdoEngine.java | 14 +++ .../ordo/api/query/InstanceQuery.java | 29 +++++ .../com/jetlumen/ordo/api/query/Page.java | 18 +++ .../jetlumen/ordo/api/query/PageRequest.java | 21 ++++ .../jetlumen/ordo/api/query/TaskQuery.java | 33 ++++++ .../repository/ApprovalTaskRepository.java | 6 + .../ProcessDefinitionRepository.java | 5 + .../repository/ProcessInstanceRepository.java | 6 + .../ordo/api/query/QueryTypesTest.java | 60 ++++++++++ .../jetlumen/ordo/core/DefaultOrdoEngine.java | 24 ++++ .../ordo/core/InMemoryOrdoEngine.java | 24 +++- .../InMemoryApprovalTaskRepository.java | 66 +++++++++++ .../InMemoryProcessDefinitionRepository.java | 16 +++ .../InMemoryProcessInstanceRepository.java | 37 ++++++ .../ordo/core/InMemoryOrdoEngineTest.java | 54 +++++++++ .../InMemoryApprovalTaskRepositoryTest.java | 81 +++++++++++++ ...MemoryProcessDefinitionRepositoryTest.java | 37 ++++++ ...InMemoryProcessInstanceRepositoryTest.java | 55 +++++++++ .../jdbc/JdbcApprovalTaskRepository.java | 71 +++++++++++ .../jdbc/JdbcProcessDefinitionRepository.java | 110 ++++++++++++------ .../jdbc/JdbcProcessInstanceRepository.java | 65 +++++++++++ .../ordo/storage/jdbc/PageSupport.java | 35 ++++++ .../jdbc/JdbcApprovalTaskRepositoryTest.java | 38 ++++++ .../JdbcProcessDefinitionRepositoryTest.java | 19 +++ .../JdbcProcessInstanceRepositoryTest.java | 38 ++++++ 26 files changed, 963 insertions(+), 39 deletions(-) create mode 100644 docs/roadmap.md create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/query/InstanceQuery.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/query/Page.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/query/PageRequest.java create mode 100644 ordo-api/src/main/java/com/jetlumen/ordo/api/query/TaskQuery.java create mode 100644 ordo-api/src/test/java/com/jetlumen/ordo/api/query/QueryTypesTest.java create mode 100644 ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java create mode 100644 ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepositoryTest.java create mode 100644 ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepositoryTest.java create mode 100644 ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/PageSupport.java diff --git a/docs/roadmap.md b/docs/roadmap.md new file mode 100644 index 0000000..3a969fa --- /dev/null +++ b/docs/roadmap.md @@ -0,0 +1,40 @@ +# Ordo Roadmap + +记录当前已完成能力之后,后续要做的开发计划。按优先级分组,供后续排期/立项参考。 + +## 已完成(背景,非本文档重点) +- 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/时间范围过滤,默认按时间降序(最新优先)。 + +## P0 — 审计与扩展点 +1. **独立历史/审计事件模型**:现在历史只能靠 `ApprovalTask.action` 字段拼凑,没有独立的流程事件表 + (谁在何时对哪个实例做了什么)。建议新增 `ProcessEvent`/`ProcessHistoryRepository`。 +2. **状态变更事件监听器**:`OrdoEngine` 目前没有任何 listener/hook,无法在任务创建、审批、实例完成时 + 被外部感知(做通知、写审计日志等)。可加 `OrdoEventListener` 扩展点,风格与 `AssigneeResolver`/ + `RoutingCondition` 一致。 +3. **ACTION 步骤执行记录持久化**:目前 `ActionHandler` 执行结果只在宿主内存里记(如 rhizome 的 + `LeaveActionHandler`),重启即丢失,且失败只打日志不影响流程状态,需要设计重试/失败处理策略。 + +## P1 — 任务生命周期完善 +4. **任务委托/转派(delegate/reassign)**:任务创建后 assignee 不可变,无法转交他人处理。 +5. **超时/升级(SLA/escalation)**:无到期时间、定时器、自动升级机制。 +6. **流程实例取消 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**:属于更大的引擎能力扩展,优先级最低,等基础能力稳定后 + 再评估是否需要。 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 5810b5d..ae5736e 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 @@ -1,5 +1,10 @@ package com.jetlumen.ordo.api; +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 java.util.List; import java.util.Optional; @@ -33,4 +38,13 @@ public interface OrdoEngine { List findTasks(String instanceId); List findPendingTasksByAssignee(String assignee); List findPendingTasksByInstanceId(String instanceId); + + /** Paginated, filterable task query; see {@link com.jetlumen.ordo.api.repository.ApprovalTaskRepository#query}. */ + Page queryTasks(TaskQuery query, PageRequest pageRequest); + + /** Paginated, filterable instance query; see {@link com.jetlumen.ordo.api.repository.ProcessInstanceRepository#query}. */ + Page queryInstances(InstanceQuery query, PageRequest pageRequest); + + /** Paginated listing of all registered process definitions. */ + Page listDefinitions(PageRequest pageRequest); } diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/query/InstanceQuery.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/InstanceQuery.java new file mode 100644 index 0000000..5754681 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/InstanceQuery.java @@ -0,0 +1,29 @@ +package com.jetlumen.ordo.api.query; + +import com.jetlumen.ordo.api.ProcessStatus; + +import java.time.Instant; + +/** Optional filters for querying process instances; a null field means "no filter on that dimension". */ +public record InstanceQuery(String definitionId, ProcessStatus status, String initiator, + Instant startedFrom, Instant startedTo) { + public static InstanceQuery any() { + return new InstanceQuery(null, null, null, null, null); + } + + public InstanceQuery withDefinitionId(String definitionId) { + return new InstanceQuery(definitionId, status, initiator, startedFrom, startedTo); + } + + public InstanceQuery withStatus(ProcessStatus status) { + return new InstanceQuery(definitionId, status, initiator, startedFrom, startedTo); + } + + public InstanceQuery withInitiator(String initiator) { + return new InstanceQuery(definitionId, status, initiator, startedFrom, startedTo); + } + + public InstanceQuery withStartedBetween(Instant startedFrom, Instant startedTo) { + return new InstanceQuery(definitionId, status, initiator, startedFrom, startedTo); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/query/Page.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/Page.java new file mode 100644 index 0000000..e41ad6b --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/Page.java @@ -0,0 +1,18 @@ +package com.jetlumen.ordo.api.query; + +import java.util.List; + +/** One page of results plus enough information to compute total pages / next-page availability. */ +public record Page(List content, long totalElements, int page, int size) { + public Page { + content = List.copyOf(content); + } + + public int totalPages() { + return size == 0 ? 0 : (int) Math.ceil((double) totalElements / size); + } + + public boolean hasNext() { + return (long) (page + 1) * size < totalElements; + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/query/PageRequest.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/PageRequest.java new file mode 100644 index 0000000..5c57207 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/PageRequest.java @@ -0,0 +1,21 @@ +package com.jetlumen.ordo.api.query; + +/** Page-number based pagination request. */ +public record PageRequest(int page, int size) { + public PageRequest { + if (page < 0) { + throw new IllegalArgumentException("page must not be negative"); + } + if (size <= 0) { + throw new IllegalArgumentException("size must be positive"); + } + } + + public static PageRequest of(int page, int size) { + return new PageRequest(page, size); + } + + public int offset() { + return page * size; + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/query/TaskQuery.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/TaskQuery.java new file mode 100644 index 0000000..c76b520 --- /dev/null +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/query/TaskQuery.java @@ -0,0 +1,33 @@ +package com.jetlumen.ordo.api.query; + +import com.jetlumen.ordo.api.TaskStatus; + +import java.time.Instant; + +/** Optional filters for querying approval tasks; a null field means "no filter on that dimension". */ +public record TaskQuery(String assignee, String instanceId, String definitionId, TaskStatus status, + Instant createdFrom, Instant createdTo) { + public static TaskQuery any() { + return new TaskQuery(null, null, null, null, null, null); + } + + public TaskQuery withAssignee(String assignee) { + return new TaskQuery(assignee, instanceId, definitionId, status, createdFrom, createdTo); + } + + public TaskQuery withInstanceId(String instanceId) { + return new TaskQuery(assignee, instanceId, definitionId, status, createdFrom, createdTo); + } + + public TaskQuery withDefinitionId(String definitionId) { + return new TaskQuery(assignee, instanceId, definitionId, status, createdFrom, createdTo); + } + + public TaskQuery withStatus(TaskStatus status) { + return new TaskQuery(assignee, instanceId, definitionId, status, createdFrom, createdTo); + } + + public TaskQuery withCreatedBetween(Instant createdFrom, Instant createdTo) { + return new TaskQuery(assignee, instanceId, definitionId, status, createdFrom, createdTo); + } +} diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java index aeab185..7729f84 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ApprovalTaskRepository.java @@ -1,6 +1,9 @@ package com.jetlumen.ordo.api.repository; import com.jetlumen.ordo.api.ApprovalTask; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import com.jetlumen.ordo.api.query.TaskQuery; import java.util.List; import java.util.Optional; @@ -27,4 +30,7 @@ public interface ApprovalTaskRepository { * @return true if the update was applied, false if the task had already been completed by a concurrent operation */ boolean completeIfPending(ApprovalTask completedTask); + + /** Paginated, filterable query; results are ordered newest-first (created_at desc). */ + Page query(TaskQuery query, PageRequest pageRequest); } diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java index dc0528b..eb967cc 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessDefinitionRepository.java @@ -1,6 +1,8 @@ package com.jetlumen.ordo.api.repository; import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; import java.util.Optional; @@ -20,4 +22,7 @@ public interface ProcessDefinitionRepository { void upsert(ProcessDefinition definition); Optional findById(String definitionId); + + /** Paginated listing of all registered definitions, ordered by id ascending. */ + Page findAll(PageRequest pageRequest); } diff --git a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java index f15a7ec..f7d8625 100644 --- a/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java +++ b/ordo-api/src/main/java/com/jetlumen/ordo/api/repository/ProcessInstanceRepository.java @@ -1,6 +1,9 @@ package com.jetlumen.ordo.api.repository; import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.query.InstanceQuery; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; import java.util.Optional; @@ -22,4 +25,7 @@ public interface ProcessInstanceRepository { * @return true if the update was applied, false if the instance was missing or already terminal */ boolean completeIfRunning(ProcessInstance completed); + + /** Paginated, filterable query; results are ordered newest-first (started_at desc). */ + Page query(InstanceQuery query, PageRequest pageRequest); } diff --git a/ordo-api/src/test/java/com/jetlumen/ordo/api/query/QueryTypesTest.java b/ordo-api/src/test/java/com/jetlumen/ordo/api/query/QueryTypesTest.java new file mode 100644 index 0000000..b0e44ab --- /dev/null +++ b/ordo-api/src/test/java/com/jetlumen/ordo/api/query/QueryTypesTest.java @@ -0,0 +1,60 @@ +package com.jetlumen.ordo.api.query; + +import com.jetlumen.ordo.api.TaskStatus; +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.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class QueryTypesTest { + @Test + void pageRequestRejectsInvalidPageOrSize() { + assertThrows(IllegalArgumentException.class, () -> new PageRequest(-1, 10)); + assertThrows(IllegalArgumentException.class, () -> new PageRequest(0, 0)); + assertEquals(20, PageRequest.of(2, 10).offset()); + } + + @Test + void pageComputesTotalPagesAndHasNext() { + Page page = new Page<>(List.of("a", "b"), 5, 0, 2); + assertEquals(3, page.totalPages()); + assertTrue(page.hasNext()); + + Page lastPage = new Page<>(List.of("e"), 5, 2, 2); + assertFalse(lastPage.hasNext()); + } + + @Test + void taskQueryAnyHasNoFiltersAndWithersAreImmutable() { + TaskQuery any = TaskQuery.any(); + assertNull(any.assignee()); + assertNull(any.definitionId()); + + TaskQuery filtered = any.withAssignee("maria").withDefinitionId("leave").withStatus(TaskStatus.PENDING); + assertEquals("maria", filtered.assignee()); + assertEquals("leave", filtered.definitionId()); + assertEquals(TaskStatus.PENDING, filtered.status()); + assertNull(any.assignee()); + } + + @Test + void instanceQueryAnyHasNoFiltersAndWithersAreImmutable() { + InstanceQuery any = InstanceQuery.any(); + assertNull(any.definitionId()); + + Instant from = Instant.parse("2026-01-01T00:00:00Z"); + Instant to = Instant.parse("2026-01-02T00:00:00Z"); + InstanceQuery filtered = any.withDefinitionId("leave").withInitiator("alice").withStartedBetween(from, to); + assertEquals("leave", filtered.definitionId()); + assertEquals("alice", filtered.initiator()); + assertEquals(from, filtered.startedFrom()); + assertEquals(to, filtered.startedTo()); + assertNull(any.definitionId()); + } +} 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 43fc3b3..07fa326 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 @@ -26,6 +26,10 @@ import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException; import com.jetlumen.ordo.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; +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.ApprovalTaskRepository; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; @@ -198,6 +202,26 @@ public final class DefaultOrdoEngine implements OrdoEngine { return taskRepository.findPendingByInstanceId(instanceId); } + @Override + public synchronized Page queryTasks(TaskQuery query, PageRequest pageRequest) { + Objects.requireNonNull(query, "query must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + return taskRepository.query(query, pageRequest); + } + + @Override + public synchronized Page queryInstances(InstanceQuery query, PageRequest pageRequest) { + Objects.requireNonNull(query, "query must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + return instanceRepository.query(query, pageRequest); + } + + @Override + public synchronized Page listDefinitions(PageRequest pageRequest) { + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + return definitionRepository.findAll(pageRequest); + } + private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now) { for (String candidate : step.candidates()) { String assignee = assigneeResolver.resolve(candidate, step, instance.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 bdb4f7f..937e74f 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 @@ -8,6 +8,10 @@ import com.jetlumen.ordo.api.ProcessContext; import com.jetlumen.ordo.api.ProcessDefinition; 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.InMemoryApprovalTaskRepository; import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository; import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository; @@ -50,11 +54,12 @@ public final class InMemoryOrdoEngine implements OrdoEngine { public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition, ActionHandler actionHandler) { + InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository(); this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, actionHandler, new NoopTransactionExecutor(), new InMemoryProcessDefinitionRepository(), - new InMemoryProcessInstanceRepository(), - new InMemoryApprovalTaskRepository()); + instanceRepository, + new InMemoryApprovalTaskRepository(instanceRepository)); } @Override @@ -111,4 +116,19 @@ public final class InMemoryOrdoEngine implements OrdoEngine { public List findPendingTasksByInstanceId(String instanceId) { return delegate.findPendingTasksByInstanceId(instanceId); } + + @Override + public Page queryTasks(TaskQuery query, PageRequest pageRequest) { + return delegate.queryTasks(query, pageRequest); + } + + @Override + public Page queryInstances(InstanceQuery query, PageRequest pageRequest) { + return delegate.queryInstances(query, pageRequest); + } + + @Override + public Page listDefinitions(PageRequest pageRequest) { + return delegate.listDefinitions(pageRequest); + } } diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java index f3a0587..2389011 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepository.java @@ -1,9 +1,15 @@ package com.jetlumen.ordo.core.repository; import com.jetlumen.ordo.api.ApprovalTask; +import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.TaskStatus; +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.ApprovalTaskRepository; +import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; +import java.util.Comparator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -12,6 +18,21 @@ import java.util.Optional; /** Development-only in-memory implementation of the task storage port. */ public final class InMemoryApprovalTaskRepository implements ApprovalTaskRepository { private final Map tasks = new LinkedHashMap<>(); + private final ProcessInstanceRepository instanceRepository; + + public InMemoryApprovalTaskRepository() { + this(null); + } + + /** + * @param instanceRepository used to resolve a task's owning instance to its definition id when + * {@link TaskQuery#definitionId()} is used as a filter; pass {@code null} + * if that filter will never be used (it then fails fast instead of silently + * being ignored). + */ + public InMemoryApprovalTaskRepository(ProcessInstanceRepository instanceRepository) { + this.instanceRepository = instanceRepository; + } @Override public synchronized void save(ApprovalTask task) { @@ -61,4 +82,49 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit tasks.put(completedTask.id(), completedTask); return true; } + + @Override + public synchronized Page query(TaskQuery query, PageRequest pageRequest) { + List matched = tasks.values().stream().filter(task -> matches(task, query)).toList(); + List sorted = matched.stream() + .sorted(Comparator.comparing(ApprovalTask::createdAt).thenComparing(ApprovalTask::id).reversed()) + .toList(); + List page = sorted.stream() + .skip((long) pageRequest.offset()) + .limit(pageRequest.size()) + .toList(); + return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size()); + } + + private boolean matches(ApprovalTask task, TaskQuery query) { + if (query.assignee() != null && !query.assignee().equals(task.assignee())) { + return false; + } + if (query.instanceId() != null && !query.instanceId().equals(task.instanceId())) { + return false; + } + if (query.status() != null && query.status() != task.status()) { + return false; + } + if (query.createdFrom() != null && task.createdAt().isBefore(query.createdFrom())) { + return false; + } + if (query.createdTo() != null && task.createdAt().isAfter(query.createdTo())) { + return false; + } + if (query.definitionId() != null && !query.definitionId().equals(definitionIdOf(task))) { + return false; + } + return true; + } + + private String definitionIdOf(ApprovalTask task) { + if (instanceRepository == null) { + throw new UnsupportedOperationException( + "TaskQuery.definitionId() filtering requires an InMemoryApprovalTaskRepository " + + "constructed with a ProcessInstanceRepository"); + } + Optional instance = instanceRepository.findById(task.instanceId()); + return instance.map(ProcessInstance::definitionId).orElse(null); + } } diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java index f4bf8c0..17a2efb 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepository.java @@ -1,9 +1,13 @@ package com.jetlumen.ordo.core.repository; import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; +import java.util.Comparator; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Optional; @@ -25,4 +29,16 @@ public final class InMemoryProcessDefinitionRepository implements ProcessDefinit public synchronized Optional findById(String definitionId) { return Optional.ofNullable(definitions.get(definitionId)); } + + @Override + public synchronized Page findAll(PageRequest pageRequest) { + List sorted = definitions.values().stream() + .sorted(Comparator.comparing(ProcessDefinition::id)) + .toList(); + List page = sorted.stream() + .skip((long) pageRequest.offset()) + .limit(pageRequest.size()) + .toList(); + return new Page<>(page, sorted.size(), pageRequest.page(), pageRequest.size()); + } } diff --git a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java index 9f7c506..608e677 100644 --- a/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java +++ b/ordo-core/src/main/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepository.java @@ -2,9 +2,14 @@ package com.jetlumen.ordo.core.repository; import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessStatus; +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.repository.ProcessInstanceRepository; +import java.util.Comparator; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Optional; @@ -43,4 +48,36 @@ public final class InMemoryProcessInstanceRepository implements ProcessInstanceR instances.put(completed.id(), completed); return true; } + + @Override + public synchronized Page query(InstanceQuery query, PageRequest pageRequest) { + List matched = instances.values().stream().filter(instance -> matches(instance, query)).toList(); + List sorted = matched.stream() + .sorted(Comparator.comparing(ProcessInstance::startedAt).thenComparing(ProcessInstance::id).reversed()) + .toList(); + List page = sorted.stream() + .skip((long) pageRequest.offset()) + .limit(pageRequest.size()) + .toList(); + return new Page<>(page, matched.size(), pageRequest.page(), pageRequest.size()); + } + + private static boolean matches(ProcessInstance instance, InstanceQuery query) { + if (query.definitionId() != null && !query.definitionId().equals(instance.definitionId())) { + return false; + } + if (query.status() != null && query.status() != instance.status()) { + return false; + } + if (query.initiator() != null && !query.initiator().equals(instance.initiator())) { + return false; + } + if (query.startedFrom() != null && instance.startedAt().isBefore(query.startedFrom())) { + return false; + } + if (query.startedTo() != null && instance.startedAt().isAfter(query.startedTo())) { + return false; + } + return true; + } } 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 da17e0d..9597cae 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 @@ -21,6 +21,10 @@ import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException; import com.jetlumen.ordo.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; +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 org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -458,4 +462,54 @@ class InMemoryOrdoEngineTest { ApprovalTask task = routingEngine.findPendingTasksByInstanceId(instance.id()).getFirst(); assertThrows(NoRouteFoundException.class, () -> routingEngine.approve(task.id(), "maria")); } + + @Test + void queryTasksFiltersByAssigneeDefinitionAndPaginates() { + engine.register(ProcessDefinition.linear("expense", "Expense request", List.of( + ApprovalStep.single("finance", "Finance approval", "frank")))); + var leaveInstance = engine.start("leave", "alice"); + engine.start("expense", "alice"); + + Page byAssignee = engine.queryTasks(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10)); + assertEquals(1, byAssignee.totalElements()); + assertEquals("maria", byAssignee.content().getFirst().assignee()); + + Page byDefinition = engine.queryTasks(TaskQuery.any().withDefinitionId("leave"), + new PageRequest(0, 10)); + assertEquals(1, byDefinition.totalElements()); + assertEquals(leaveInstance.id(), byDefinition.content().getFirst().instanceId()); + + Page firstPage = engine.queryTasks(TaskQuery.any(), new PageRequest(0, 1)); + assertEquals(2, firstPage.totalElements()); + assertEquals(2, firstPage.totalPages()); + assertEquals(1, firstPage.content().size()); + assertTrue(firstPage.hasNext()); + } + + @Test + void queryInstancesFiltersByDefinitionAndStatus() { + var running = engine.start("leave", "alice"); + var toWithdraw = engine.start("leave", "bob"); + engine.withdraw(toWithdraw.id(), "bob"); + + Page runningOnly = engine.queryInstances( + InstanceQuery.any().withStatus(ProcessStatus.RUNNING), new PageRequest(0, 10)); + assertEquals(1, runningOnly.totalElements()); + assertEquals(running.id(), runningOnly.content().getFirst().id()); + + Page byInitiator = engine.queryInstances(InstanceQuery.any().withInitiator("bob"), + new PageRequest(0, 10)); + assertEquals(1, byInitiator.totalElements()); + assertEquals(ProcessStatus.WITHDRAWN, byInitiator.content().getFirst().status()); + } + + @Test + void listDefinitionsPaginatesRegisteredDefinitions() { + engine.register(ProcessDefinition.linear("expense", "Expense request", List.of( + ApprovalStep.single("finance", "Finance approval", "frank")))); + + Page all = engine.listDefinitions(new PageRequest(0, 10)); + assertEquals(2, all.totalElements()); + assertEquals(List.of("expense", "leave"), all.content().stream().map(ProcessDefinition::id).toList()); + } } diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java new file mode 100644 index 0000000..107f23a --- /dev/null +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryApprovalTaskRepositoryTest.java @@ -0,0 +1,81 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ApprovalTask; +import com.jetlumen.ordo.api.ProcessContext; +import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.ProcessStatus; +import com.jetlumen.ordo.api.TaskStatus; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import com.jetlumen.ordo.api.query.TaskQuery; +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.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class InMemoryApprovalTaskRepositoryTest { + private static final Instant CREATED_AT = Instant.parse("2026-01-15T09:00:00Z"); + + @Test + void queryFiltersPaginatesAndOrdersNewestFirst() { + InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository(); + ApprovalTask t1 = task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT); + ApprovalTask t2 = task("task-2", "inst-1", "henry", TaskStatus.PENDING, CREATED_AT.plusSeconds(5)); + ApprovalTask t3 = task("task-3", "inst-2", "maria", TaskStatus.APPROVED, CREATED_AT.plusSeconds(10)); + repository.save(t1); + repository.save(t2); + repository.save(t3); + + Page byAssignee = repository.query(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10)); + assertEquals(2, byAssignee.totalElements()); + assertEquals(List.of(t3, t1), byAssignee.content()); + + Page byStatus = repository.query(TaskQuery.any().withStatus(TaskStatus.PENDING), new PageRequest(0, 10)); + assertEquals(List.of(t2, t1), byStatus.content()); + + Page pageOne = repository.query(TaskQuery.any(), new PageRequest(0, 2)); + assertEquals(3, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of(t3, t2), pageOne.content()); + + Page pageTwo = repository.query(TaskQuery.any(), new PageRequest(1, 2)); + assertEquals(List.of(t1), pageTwo.content()); + assertFalse(pageTwo.hasNext()); + } + + @Test + void definitionIdFilterFailsFastWithoutAnInstanceRepository() { + InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository(); + repository.save(task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT)); + + assertThrows(UnsupportedOperationException.class, + () -> repository.query(TaskQuery.any().withDefinitionId("leave"), new PageRequest(0, 10))); + } + + @Test + void definitionIdFilterResolvesThroughTheInstanceRepository() { + InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository(); + instanceRepository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, + CREATED_AT, null, ProcessContext.empty())); + instanceRepository.insert(new ProcessInstance("inst-2", "expense", "alice", ProcessStatus.RUNNING, + CREATED_AT, null, ProcessContext.empty())); + InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository(instanceRepository); + ApprovalTask leaveTask = task("task-1", "inst-1", "maria", TaskStatus.PENDING, CREATED_AT); + ApprovalTask expenseTask = task("task-2", "inst-2", "frank", TaskStatus.PENDING, CREATED_AT); + repository.save(leaveTask); + repository.save(expenseTask); + + Page byDefinition = repository.query(TaskQuery.any().withDefinitionId("leave"), new PageRequest(0, 10)); + assertEquals(List.of(leaveTask), byDefinition.content()); + } + + private static ApprovalTask task(String id, String instanceId, String assignee, TaskStatus status, Instant createdAt) { + return new ApprovalTask(id, instanceId, "step", "Step", assignee, status, createdAt, null, null); + } +} diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepositoryTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepositoryTest.java new file mode 100644 index 0000000..fc25cdf --- /dev/null +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessDefinitionRepositoryTest.java @@ -0,0 +1,37 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ApprovalStep; +import com.jetlumen.ordo.api.ProcessDefinition; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import org.junit.jupiter.api.Test; + +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 InMemoryProcessDefinitionRepositoryTest { + @Test + void findAllPaginatesDefinitionsOrderedById() { + InMemoryProcessDefinitionRepository repository = new InMemoryProcessDefinitionRepository(); + repository.insertIfAbsent(definition("c-def")); + repository.insertIfAbsent(definition("a-def")); + repository.insertIfAbsent(definition("b-def")); + + Page pageOne = repository.findAll(new PageRequest(0, 2)); + assertEquals(3, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of("a-def", "b-def"), pageOne.content().stream().map(ProcessDefinition::id).toList()); + + Page pageTwo = repository.findAll(new PageRequest(1, 2)); + assertEquals(List.of("c-def"), pageTwo.content().stream().map(ProcessDefinition::id).toList()); + assertFalse(pageTwo.hasNext()); + } + + private static ProcessDefinition definition(String id) { + return ProcessDefinition.linear(id, id, List.of(ApprovalStep.single("lead", "Lead approval", "lee"))); + } +} diff --git a/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepositoryTest.java b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepositoryTest.java new file mode 100644 index 0000000..03fde1a --- /dev/null +++ b/ordo-core/src/test/java/com/jetlumen/ordo/core/repository/InMemoryProcessInstanceRepositoryTest.java @@ -0,0 +1,55 @@ +package com.jetlumen.ordo.core.repository; + +import com.jetlumen.ordo.api.ProcessContext; +import com.jetlumen.ordo.api.ProcessInstance; +import com.jetlumen.ordo.api.ProcessStatus; +import com.jetlumen.ordo.api.query.InstanceQuery; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +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 InMemoryProcessInstanceRepositoryTest { + private static final Instant STARTED_AT = Instant.parse("2026-01-15T09:00:00Z"); + + @Test + void queryFiltersPaginatesAndOrdersNewestFirst() { + InMemoryProcessInstanceRepository repository = new InMemoryProcessInstanceRepository(); + ProcessInstance i1 = instance("inst-1", "leave", "alice", ProcessStatus.RUNNING, STARTED_AT); + ProcessInstance i2 = instance("inst-2", "leave", "bob", ProcessStatus.RUNNING, STARTED_AT.plusSeconds(5)); + ProcessInstance i3 = instance("inst-3", "expense", "alice", ProcessStatus.WITHDRAWN, STARTED_AT.plusSeconds(10)); + repository.insert(i1); + repository.insert(i2); + repository.insert(i3); + + Page byDefinition = repository.query(InstanceQuery.any().withDefinitionId("leave"), + new PageRequest(0, 10)); + assertEquals(2, byDefinition.totalElements()); + assertEquals(List.of(i2, i1), byDefinition.content()); + + Page byStatus = repository.query(InstanceQuery.any().withStatus(ProcessStatus.WITHDRAWN), + new PageRequest(0, 10)); + assertEquals(List.of(i3), byStatus.content()); + + Page pageOne = repository.query(InstanceQuery.any(), new PageRequest(0, 2)); + assertEquals(3, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of(i3, i2), pageOne.content()); + + Page pageTwo = repository.query(InstanceQuery.any(), new PageRequest(1, 2)); + assertEquals(List.of(i1), pageTwo.content()); + assertFalse(pageTwo.hasNext()); + } + + private static ProcessInstance instance(String id, String definitionId, String initiator, ProcessStatus status, + Instant startedAt) { + return new ProcessInstance(id, definitionId, initiator, status, startedAt, null, ProcessContext.empty()); + } +} diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java index 770e061..0103f01 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepository.java @@ -1,13 +1,18 @@ package com.jetlumen.ordo.storage.jdbc; import com.jetlumen.ordo.api.ApprovalTask; +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.ApprovalTaskRepository; +import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; +import java.sql.Timestamp; import java.util.ArrayList; import java.util.List; import java.util.Objects; @@ -36,6 +41,9 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository private static final String SELECT_PENDING_BY_INSTANCE = "SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND instance_id = ?" + " ORDER BY created_at, id"; + private static final String TASK_COLUMNS_QUALIFIED = + "t.id, t.instance_id, t.step_id, t.task_name, t.assignee, t.status, t.created_at, t.completed_at," + + " t.action_actor, t.action_comment, t.action_at"; private final JdbcConnectionProvider connectionProvider; @@ -111,6 +119,69 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository } } + @Override + public Page query(TaskQuery query, PageRequest pageRequest) { + Objects.requireNonNull(query, "query must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + WhereClause where = buildWhereClause(query); + Connection connection = connectionProvider.getConnection(); + try { + long total = PageSupport.count(connection, + "SELECT COUNT(*) FROM ordo_approval_task t" + where.sql(), where.params()); + String sql = "SELECT " + TASK_COLUMNS_QUALIFIED + " FROM ordo_approval_task t" + where.sql() + + " ORDER BY t.created_at DESC, t.id DESC LIMIT ? OFFSET ?"; + List content = new ArrayList<>(); + try (PreparedStatement select = connection.prepareStatement(sql)) { + int index = PageSupport.bindParams(select, where.params()); + select.setInt(index++, pageRequest.size()); + select.setInt(index, pageRequest.offset()); + try (ResultSet resultSet = select.executeQuery()) { + while (resultSet.next()) { + content.add(ApprovalTaskMapper.read(resultSet)); + } + } + } + return new Page<>(content, total, pageRequest.page(), pageRequest.size()); + } catch (SQLException e) { + throw new JdbcStorageException("failed to query tasks", e); + } finally { + connectionProvider.close(connection); + } + } + + private static WhereClause buildWhereClause(TaskQuery query) { + List conditions = new ArrayList<>(); + List params = new ArrayList<>(); + if (query.assignee() != null) { + conditions.add("t.assignee = ?"); + params.add(query.assignee()); + } + if (query.instanceId() != null) { + conditions.add("t.instance_id = ?"); + params.add(query.instanceId()); + } + if (query.status() != null) { + conditions.add("t.status = ?"); + params.add(query.status().name()); + } + if (query.createdFrom() != null) { + conditions.add("t.created_at >= ?"); + params.add(Timestamp.from(query.createdFrom())); + } + if (query.createdTo() != null) { + conditions.add("t.created_at <= ?"); + params.add(Timestamp.from(query.createdTo())); + } + String join = ""; + if (query.definitionId() != null) { + join = " JOIN ordo_process_instance i ON i.id = t.instance_id"; + conditions.add("i.definition_id = ?"); + params.add(query.definitionId()); + } + String where = conditions.isEmpty() ? "" : " WHERE " + String.join(" AND ", conditions); + return new WhereClause(join + where, params); + } + private Optional findOne(String sql, String parameter) { Connection connection = connectionProvider.getConnection(); try (PreparedStatement select = connection.prepareStatement(sql)) { diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java index c41c7ad..d690ed9 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepository.java @@ -3,6 +3,8 @@ package com.jetlumen.ordo.storage.jdbc; import com.jetlumen.ordo.api.ApprovalStep; import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.StepTransition; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper.StepRow; @@ -49,6 +51,8 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR private static final String SELECT_TRANSITIONS = "SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition " + "WHERE definition_id = ? ORDER BY from_step_id, priority"; + private static final String SELECT_DEFINITIONS_PAGE = + "SELECT id, name FROM ordo_process_definition ORDER BY id LIMIT ? OFFSET ?"; private final JdbcConnectionProvider connectionProvider; @@ -112,43 +116,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR name = resultSet.getString("name"); } } - List stepRows = new ArrayList<>(); - try (PreparedStatement selectSteps = connection.prepareStatement(SELECT_STEPS)) { - selectSteps.setString(1, definitionId); - try (ResultSet resultSet = selectSteps.executeQuery()) { - while (resultSet.next()) { - stepRows.add(ApprovalStepMapper.readRow(resultSet)); - } - } - } - Map> candidatesByStep = new LinkedHashMap<>(); - try (PreparedStatement selectCandidates = connection.prepareStatement(SELECT_CANDIDATES)) { - selectCandidates.setString(1, definitionId); - try (ResultSet resultSet = selectCandidates.executeQuery()) { - while (resultSet.next()) { - candidatesByStep.computeIfAbsent(resultSet.getString("step_id"), key -> new ArrayList<>()) - .add(resultSet.getString("candidate")); - } - } - } - List steps = new ArrayList<>(); - for (StepRow row : stepRows) { - List candidates = candidatesByStep.getOrDefault(row.stepId(), List.of()); - steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(), - row.actionKey())); - } - List transitions = new ArrayList<>(); - try (PreparedStatement selectTransitions = connection.prepareStatement(SELECT_TRANSITIONS)) { - selectTransitions.setString(1, definitionId); - try (ResultSet resultSet = selectTransitions.executeQuery()) { - while (resultSet.next()) { - TransitionRow row = StepTransitionMapper.readRow(resultSet); - transitions.add(new StepTransition(row.fromStepId(), row.toStepId(), row.conditionKey(), - row.priority())); - } - } - } - return Optional.of(new ProcessDefinition(definitionId, name, steps, transitions)); + return Optional.of(assemble(connection, definitionId, name)); } catch (SQLException e) { throw new JdbcStorageException("failed to load definition: " + definitionId, e); } finally { @@ -156,6 +124,74 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR } } + @Override + public Page findAll(PageRequest pageRequest) { + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + Connection connection = connectionProvider.getConnection(); + try { + long total = PageSupport.count(connection, "SELECT COUNT(*) FROM ordo_process_definition", List.of()); + List content = new ArrayList<>(); + try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITIONS_PAGE)) { + select.setInt(1, pageRequest.size()); + select.setInt(2, pageRequest.offset()); + List> idsAndNames = new ArrayList<>(); + try (ResultSet resultSet = select.executeQuery()) { + while (resultSet.next()) { + idsAndNames.add(Map.entry(resultSet.getString("id"), resultSet.getString("name"))); + } + } + for (Map.Entry idAndName : idsAndNames) { + content.add(assemble(connection, idAndName.getKey(), idAndName.getValue())); + } + } + return new Page<>(content, total, pageRequest.page(), pageRequest.size()); + } catch (SQLException e) { + throw new JdbcStorageException("failed to list definitions", e); + } finally { + connectionProvider.close(connection); + } + } + + private static ProcessDefinition assemble(Connection connection, String definitionId, String name) throws SQLException { + List stepRows = new ArrayList<>(); + try (PreparedStatement selectSteps = connection.prepareStatement(SELECT_STEPS)) { + selectSteps.setString(1, definitionId); + try (ResultSet resultSet = selectSteps.executeQuery()) { + while (resultSet.next()) { + stepRows.add(ApprovalStepMapper.readRow(resultSet)); + } + } + } + Map> candidatesByStep = new LinkedHashMap<>(); + try (PreparedStatement selectCandidates = connection.prepareStatement(SELECT_CANDIDATES)) { + selectCandidates.setString(1, definitionId); + try (ResultSet resultSet = selectCandidates.executeQuery()) { + while (resultSet.next()) { + candidatesByStep.computeIfAbsent(resultSet.getString("step_id"), key -> new ArrayList<>()) + .add(resultSet.getString("candidate")); + } + } + } + List steps = new ArrayList<>(); + for (StepRow row : stepRows) { + List candidates = candidatesByStep.getOrDefault(row.stepId(), List.of()); + steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(), + row.actionKey())); + } + List transitions = new ArrayList<>(); + try (PreparedStatement selectTransitions = connection.prepareStatement(SELECT_TRANSITIONS)) { + selectTransitions.setString(1, definitionId); + try (ResultSet resultSet = selectTransitions.executeQuery()) { + while (resultSet.next()) { + TransitionRow row = StepTransitionMapper.readRow(resultSet); + transitions.add(new StepTransition(row.fromStepId(), row.toStepId(), row.conditionKey(), + row.priority())); + } + } + } + return new ProcessDefinition(definitionId, name, steps, transitions); + } + private static boolean definitionExists(Connection connection, String definitionId) throws SQLException { try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITION)) { select.setString(1, definitionId); diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java index b6c56b7..b6f49b0 100644 --- a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepository.java @@ -2,13 +2,20 @@ package com.jetlumen.ordo.storage.jdbc; import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessStatus; +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.repository.ProcessInstanceRepository; +import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause; import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; +import java.sql.Timestamp; +import java.util.ArrayList; +import java.util.List; import java.util.Objects; import java.util.Optional; @@ -107,4 +114,62 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos connectionProvider.close(connection); } } + + @Override + public Page query(InstanceQuery query, PageRequest pageRequest) { + Objects.requireNonNull(query, "query must not be null"); + Objects.requireNonNull(pageRequest, "pageRequest must not be null"); + WhereClause where = buildWhereClause(query); + Connection connection = connectionProvider.getConnection(); + try { + long total = PageSupport.count(connection, + "SELECT COUNT(*) FROM ordo_process_instance" + where.sql(), where.params()); + String sql = "SELECT id, definition_id, initiator, status, context_json, started_at, finished_at" + + " FROM ordo_process_instance" + where.sql() + + " ORDER BY started_at DESC, id DESC LIMIT ? OFFSET ?"; + List content = new ArrayList<>(); + try (PreparedStatement select = connection.prepareStatement(sql)) { + int index = PageSupport.bindParams(select, where.params()); + select.setInt(index++, pageRequest.size()); + select.setInt(index, pageRequest.offset()); + try (ResultSet resultSet = select.executeQuery()) { + while (resultSet.next()) { + content.add(ProcessInstanceMapper.read(resultSet)); + } + } + } + return new Page<>(content, total, pageRequest.page(), pageRequest.size()); + } catch (SQLException e) { + throw new JdbcStorageException("failed to query instances", e); + } finally { + connectionProvider.close(connection); + } + } + + private static WhereClause buildWhereClause(InstanceQuery query) { + List conditions = new ArrayList<>(); + List params = new ArrayList<>(); + if (query.definitionId() != null) { + conditions.add("definition_id = ?"); + params.add(query.definitionId()); + } + if (query.status() != null) { + conditions.add("status = ?"); + params.add(query.status().name()); + } + if (query.initiator() != null) { + conditions.add("initiator = ?"); + params.add(query.initiator()); + } + if (query.startedFrom() != null) { + conditions.add("started_at >= ?"); + params.add(Timestamp.from(query.startedFrom())); + } + if (query.startedTo() != null) { + conditions.add("started_at <= ?"); + params.add(Timestamp.from(query.startedTo())); + } + String where = conditions.isEmpty() ? "" : " WHERE " + String.join(" AND ", conditions); + return new WhereClause(where, params); + } } diff --git a/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/PageSupport.java b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/PageSupport.java new file mode 100644 index 0000000..888ae80 --- /dev/null +++ b/ordo-storage-jdbc/src/main/java/com/jetlumen/ordo/storage/jdbc/PageSupport.java @@ -0,0 +1,35 @@ +package com.jetlumen.ordo.storage.jdbc; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.List; + +/** Shared helpers for building dynamic, paginated {@code WHERE} clauses across JDBC repositories. */ +final class PageSupport { + private PageSupport() { + } + + /** A dynamically built {@code WHERE}/{@code JOIN} fragment (may be empty) plus its bind parameters, in order. */ + record WhereClause(String sql, List params) { + } + + static long count(Connection connection, String countSql, List params) throws SQLException { + try (PreparedStatement select = connection.prepareStatement(countSql)) { + bindParams(select, params); + try (ResultSet resultSet = select.executeQuery()) { + resultSet.next(); + return resultSet.getLong(1); + } + } + } + + static int bindParams(PreparedStatement statement, List params) throws SQLException { + int index = 1; + for (Object param : params) { + statement.setObject(index++, param); + } + return index; + } +} diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java index f6c98d3..a024bf4 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcApprovalTaskRepositoryTest.java @@ -8,6 +8,9 @@ import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessStatus; import com.jetlumen.ordo.api.TaskAction; import com.jetlumen.ordo.api.TaskStatus; +import com.jetlumen.ordo.api.query.Page; +import com.jetlumen.ordo.api.query.PageRequest; +import com.jetlumen.ordo.api.query.TaskQuery; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -95,6 +98,41 @@ class JdbcApprovalTaskRepositoryTest { assertEquals(List.of(henryTask), repository.findPendingByAssignee("henry")); } + @Test + void queryFiltersPaginatesAndOrdersNewestFirst() { + new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("expense", + "Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank")))); + JdbcProcessInstanceRepository instanceRepository = new JdbcProcessInstanceRepository(connectionProvider); + instanceRepository.insert(new ProcessInstance("inst-3", "expense", "alice", ProcessStatus.RUNNING, + CREATED_AT, null, ProcessContext.empty())); + + ApprovalTask t1 = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT); + ApprovalTask t2 = pendingTask("task-2", "inst-1", "hr", "henry", CREATED_AT.plusSeconds(5)); + ApprovalTask t3 = pendingTask("task-3", "inst-2", "manager", "maria", CREATED_AT.plusSeconds(10)); + ApprovalTask t4 = pendingTask("task-4", "inst-3", "finance", "frank", CREATED_AT.plusSeconds(15)); + repository.save(t1); + repository.save(t2); + repository.save(t3); + repository.save(t4); + + Page byAssignee = repository.query(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10)); + assertEquals(2, byAssignee.totalElements()); + assertEquals(List.of(t3, t1), byAssignee.content()); + + Page byDefinition = repository.query(TaskQuery.any().withDefinitionId("leave"), new PageRequest(0, 10)); + assertEquals(3, byDefinition.totalElements()); + + Page pageOne = repository.query(TaskQuery.any(), new PageRequest(0, 2)); + assertEquals(4, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of(t4, t3), pageOne.content()); + + Page pageTwo = repository.query(TaskQuery.any(), new PageRequest(1, 2)); + assertEquals(List.of(t2, t1), pageTwo.content()); + assertFalse(pageTwo.hasNext()); + } + @Test void onlyOneOfTwoConcurrentCompletionsWins() throws Exception { ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT); diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java index d047382..f1f9ef8 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessDefinitionRepositoryTest.java @@ -4,6 +4,8 @@ import com.jetlumen.ordo.api.ApprovalStep; import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.StepKind; import com.jetlumen.ordo.api.StepTransition; +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; @@ -91,4 +93,21 @@ class JdbcProcessDefinitionRepositoryTest { private static ProcessDefinition definition(String id, String name) { return ProcessDefinition.linear(id, name, List.of(ApprovalStep.single("lead", "Lead approval", "lee"))); } + + @Test + void findAllPaginatesDefinitionsOrderedById() { + repository.insertIfAbsent(definition("c-def", "C")); + repository.insertIfAbsent(definition("a-def", "A")); + repository.insertIfAbsent(definition("b-def", "B")); + + Page pageOne = repository.findAll(new PageRequest(0, 2)); + assertEquals(3, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of("a-def", "b-def"), pageOne.content().stream().map(ProcessDefinition::id).toList()); + + Page pageTwo = repository.findAll(new PageRequest(1, 2)); + assertEquals(List.of("c-def"), pageTwo.content().stream().map(ProcessDefinition::id).toList()); + assertFalse(pageTwo.hasNext()); + } } diff --git a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java index 80cd099..088fa48 100644 --- a/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java +++ b/ordo-storage-jdbc/src/test/java/com/jetlumen/ordo/storage/jdbc/JdbcProcessInstanceRepositoryTest.java @@ -5,6 +5,9 @@ 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.InstanceQuery; +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; @@ -105,4 +108,39 @@ class JdbcProcessInstanceRepositoryTest { assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave", "alice", ProcessStatus.APPROVED, STARTED_AT, STARTED_AT, ProcessContext.empty()))); } + + @Test + void queryFiltersPaginatesAndOrdersNewestFirst() { + new JdbcProcessDefinitionRepository(connectionProvider).insertIfAbsent(ProcessDefinition.linear("expense", + "Expense request", List.of(ApprovalStep.single("finance", "Finance approval", "frank")))); + + ProcessInstance i1 = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING, + STARTED_AT, null, ProcessContext.empty()); + ProcessInstance i2 = new ProcessInstance("inst-2", "leave", "bob", ProcessStatus.RUNNING, + STARTED_AT.plusSeconds(5), null, ProcessContext.empty()); + ProcessInstance i3 = new ProcessInstance("inst-3", "expense", "alice", ProcessStatus.RUNNING, + STARTED_AT.plusSeconds(10), null, ProcessContext.empty()); + repository.insert(i1); + repository.insert(i2); + repository.insert(i3); + + Page byDefinition = repository.query(InstanceQuery.any().withDefinitionId("leave"), + new PageRequest(0, 10)); + assertEquals(2, byDefinition.totalElements()); + assertEquals(List.of(i2, i1), byDefinition.content()); + + Page byInitiator = repository.query(InstanceQuery.any().withInitiator("alice"), + new PageRequest(0, 10)); + assertEquals(List.of(i3, i1), byInitiator.content()); + + Page pageOne = repository.query(InstanceQuery.any(), new PageRequest(0, 2)); + assertEquals(3, pageOne.totalElements()); + assertEquals(2, pageOne.totalPages()); + assertTrue(pageOne.hasNext()); + assertEquals(List.of(i3, i2), pageOne.content()); + + Page pageTwo = repository.query(InstanceQuery.any(), new PageRequest(1, 2)); + assertEquals(List.of(i1), pageTwo.content()); + assertFalse(pageTwo.hasNext()); + } }