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.
This commit is contained in:
0264408
2026-09-15 08:56:05 +08:00
parent e83234e98a
commit 77d1e5198b
26 changed files with 963 additions and 39 deletions
+40
View File
@@ -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**:属于更大的引擎能力扩展,优先级最低,等基础能力稳定后
再评估是否需要。
@@ -1,5 +1,10 @@
package com.jetlumen.ordo.api; 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.List;
import java.util.Optional; import java.util.Optional;
@@ -33,4 +38,13 @@ public interface OrdoEngine {
List<ApprovalTask> findTasks(String instanceId); List<ApprovalTask> findTasks(String instanceId);
List<ApprovalTask> findPendingTasksByAssignee(String assignee); List<ApprovalTask> findPendingTasksByAssignee(String assignee);
List<ApprovalTask> findPendingTasksByInstanceId(String instanceId); List<ApprovalTask> findPendingTasksByInstanceId(String instanceId);
/** Paginated, filterable task query; see {@link com.jetlumen.ordo.api.repository.ApprovalTaskRepository#query}. */
Page<ApprovalTask> queryTasks(TaskQuery query, PageRequest pageRequest);
/** Paginated, filterable instance query; see {@link com.jetlumen.ordo.api.repository.ProcessInstanceRepository#query}. */
Page<ProcessInstance> queryInstances(InstanceQuery query, PageRequest pageRequest);
/** Paginated listing of all registered process definitions. */
Page<ProcessDefinition> listDefinitions(PageRequest pageRequest);
} }
@@ -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);
}
}
@@ -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<T>(List<T> 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;
}
}
@@ -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;
}
}
@@ -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);
}
}
@@ -1,6 +1,9 @@
package com.jetlumen.ordo.api.repository; package com.jetlumen.ordo.api.repository;
import com.jetlumen.ordo.api.ApprovalTask; 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.List;
import java.util.Optional; 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 * @return true if the update was applied, false if the task had already been completed by a concurrent operation
*/ */
boolean completeIfPending(ApprovalTask completedTask); boolean completeIfPending(ApprovalTask completedTask);
/** Paginated, filterable query; results are ordered newest-first (created_at desc). */
Page<ApprovalTask> query(TaskQuery query, PageRequest pageRequest);
} }
@@ -1,6 +1,8 @@
package com.jetlumen.ordo.api.repository; package com.jetlumen.ordo.api.repository;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.query.Page;
import com.jetlumen.ordo.api.query.PageRequest;
import java.util.Optional; import java.util.Optional;
@@ -20,4 +22,7 @@ public interface ProcessDefinitionRepository {
void upsert(ProcessDefinition definition); void upsert(ProcessDefinition definition);
Optional<ProcessDefinition> findById(String definitionId); Optional<ProcessDefinition> findById(String definitionId);
/** Paginated listing of all registered definitions, ordered by id ascending. */
Page<ProcessDefinition> findAll(PageRequest pageRequest);
} }
@@ -1,6 +1,9 @@
package com.jetlumen.ordo.api.repository; package com.jetlumen.ordo.api.repository;
import com.jetlumen.ordo.api.ProcessInstance; 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; 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 * @return true if the update was applied, false if the instance was missing or already terminal
*/ */
boolean completeIfRunning(ProcessInstance completed); boolean completeIfRunning(ProcessInstance completed);
/** Paginated, filterable query; results are ordered newest-first (started_at desc). */
Page<ProcessInstance> query(InstanceQuery query, PageRequest pageRequest);
} }
@@ -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<String> page = new Page<>(List.of("a", "b"), 5, 0, 2);
assertEquals(3, page.totalPages());
assertTrue(page.hasNext());
Page<String> 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());
}
}
@@ -26,6 +26,10 @@ import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.TaskNotFoundException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; 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.ApprovalTaskRepository;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository; import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
@@ -198,6 +202,26 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return taskRepository.findPendingByInstanceId(instanceId); return taskRepository.findPendingByInstanceId(instanceId);
} }
@Override
public synchronized Page<ApprovalTask> 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<ProcessInstance> 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<ProcessDefinition> listDefinitions(PageRequest pageRequest) {
Objects.requireNonNull(pageRequest, "pageRequest must not be null");
return definitionRepository.findAll(pageRequest);
}
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now) { private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now) {
for (String candidate : step.candidates()) { for (String candidate : step.candidates()) {
String assignee = assigneeResolver.resolve(candidate, step, instance.context()); String assignee = assigneeResolver.resolve(candidate, step, instance.context());
@@ -8,6 +8,10 @@ import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.RoutingCondition; 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.InMemoryApprovalTaskRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository; import com.jetlumen.ordo.core.repository.InMemoryProcessDefinitionRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository; 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, public InMemoryOrdoEngine(Clock clock, AssigneeResolver assigneeResolver, RoutingCondition routingCondition,
ActionHandler actionHandler) { ActionHandler actionHandler) {
InMemoryProcessInstanceRepository instanceRepository = new InMemoryProcessInstanceRepository();
this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, actionHandler, this.delegate = new DefaultOrdoEngine(clock, assigneeResolver, routingCondition, actionHandler,
new NoopTransactionExecutor(), new NoopTransactionExecutor(),
new InMemoryProcessDefinitionRepository(), new InMemoryProcessDefinitionRepository(),
new InMemoryProcessInstanceRepository(), instanceRepository,
new InMemoryApprovalTaskRepository()); new InMemoryApprovalTaskRepository(instanceRepository));
} }
@Override @Override
@@ -111,4 +116,19 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
public List<ApprovalTask> findPendingTasksByInstanceId(String instanceId) { public List<ApprovalTask> findPendingTasksByInstanceId(String instanceId) {
return delegate.findPendingTasksByInstanceId(instanceId); return delegate.findPendingTasksByInstanceId(instanceId);
} }
@Override
public Page<ApprovalTask> queryTasks(TaskQuery query, PageRequest pageRequest) {
return delegate.queryTasks(query, pageRequest);
}
@Override
public Page<ProcessInstance> queryInstances(InstanceQuery query, PageRequest pageRequest) {
return delegate.queryInstances(query, pageRequest);
}
@Override
public Page<ProcessDefinition> listDefinitions(PageRequest pageRequest) {
return delegate.listDefinitions(pageRequest);
}
} }
@@ -1,9 +1,15 @@
package com.jetlumen.ordo.core.repository; package com.jetlumen.ordo.core.repository;
import com.jetlumen.ordo.api.ApprovalTask; import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.TaskStatus; 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.ApprovalTaskRepository;
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
import java.util.Comparator;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
@@ -12,6 +18,21 @@ import java.util.Optional;
/** Development-only in-memory implementation of the task storage port. */ /** Development-only in-memory implementation of the task storage port. */
public final class InMemoryApprovalTaskRepository implements ApprovalTaskRepository { public final class InMemoryApprovalTaskRepository implements ApprovalTaskRepository {
private final Map<String, ApprovalTask> tasks = new LinkedHashMap<>(); private final Map<String, ApprovalTask> 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 @Override
public synchronized void save(ApprovalTask task) { public synchronized void save(ApprovalTask task) {
@@ -61,4 +82,49 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
tasks.put(completedTask.id(), completedTask); tasks.put(completedTask.id(), completedTask);
return true; return true;
} }
@Override
public synchronized Page<ApprovalTask> query(TaskQuery query, PageRequest pageRequest) {
List<ApprovalTask> matched = tasks.values().stream().filter(task -> matches(task, query)).toList();
List<ApprovalTask> sorted = matched.stream()
.sorted(Comparator.comparing(ApprovalTask::createdAt).thenComparing(ApprovalTask::id).reversed())
.toList();
List<ApprovalTask> 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<ProcessInstance> instance = instanceRepository.findById(task.instanceId());
return instance.map(ProcessInstance::definitionId).orElse(null);
}
} }
@@ -1,9 +1,13 @@
package com.jetlumen.ordo.core.repository; package com.jetlumen.ordo.core.repository;
import com.jetlumen.ordo.api.ProcessDefinition; 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 com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
import java.util.Comparator;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
@@ -25,4 +29,16 @@ public final class InMemoryProcessDefinitionRepository implements ProcessDefinit
public synchronized Optional<ProcessDefinition> findById(String definitionId) { public synchronized Optional<ProcessDefinition> findById(String definitionId) {
return Optional.ofNullable(definitions.get(definitionId)); return Optional.ofNullable(definitions.get(definitionId));
} }
@Override
public synchronized Page<ProcessDefinition> findAll(PageRequest pageRequest) {
List<ProcessDefinition> sorted = definitions.values().stream()
.sorted(Comparator.comparing(ProcessDefinition::id))
.toList();
List<ProcessDefinition> page = sorted.stream()
.skip((long) pageRequest.offset())
.limit(pageRequest.size())
.toList();
return new Page<>(page, sorted.size(), pageRequest.page(), pageRequest.size());
}
} }
@@ -2,9 +2,14 @@ package com.jetlumen.ordo.core.repository;
import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus; 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.api.repository.ProcessInstanceRepository;
import java.util.Comparator;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
@@ -43,4 +48,36 @@ public final class InMemoryProcessInstanceRepository implements ProcessInstanceR
instances.put(completed.id(), completed); instances.put(completed.id(), completed);
return true; return true;
} }
@Override
public synchronized Page<ProcessInstance> query(InstanceQuery query, PageRequest pageRequest) {
List<ProcessInstance> matched = instances.values().stream().filter(instance -> matches(instance, query)).toList();
List<ProcessInstance> sorted = matched.stream()
.sorted(Comparator.comparing(ProcessInstance::startedAt).thenComparing(ProcessInstance::id).reversed())
.toList();
List<ProcessInstance> 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;
}
} }
@@ -21,6 +21,10 @@ import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.TaskNotFoundException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -458,4 +462,54 @@ class InMemoryOrdoEngineTest {
ApprovalTask task = routingEngine.findPendingTasksByInstanceId(instance.id()).getFirst(); ApprovalTask task = routingEngine.findPendingTasksByInstanceId(instance.id()).getFirst();
assertThrows(NoRouteFoundException.class, () -> routingEngine.approve(task.id(), "maria")); 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<ApprovalTask> byAssignee = engine.queryTasks(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10));
assertEquals(1, byAssignee.totalElements());
assertEquals("maria", byAssignee.content().getFirst().assignee());
Page<ApprovalTask> byDefinition = engine.queryTasks(TaskQuery.any().withDefinitionId("leave"),
new PageRequest(0, 10));
assertEquals(1, byDefinition.totalElements());
assertEquals(leaveInstance.id(), byDefinition.content().getFirst().instanceId());
Page<ApprovalTask> 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<ProcessInstance> runningOnly = engine.queryInstances(
InstanceQuery.any().withStatus(ProcessStatus.RUNNING), new PageRequest(0, 10));
assertEquals(1, runningOnly.totalElements());
assertEquals(running.id(), runningOnly.content().getFirst().id());
Page<ProcessInstance> 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<ProcessDefinition> all = engine.listDefinitions(new PageRequest(0, 10));
assertEquals(2, all.totalElements());
assertEquals(List.of("expense", "leave"), all.content().stream().map(ProcessDefinition::id).toList());
}
} }
@@ -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<ApprovalTask> byAssignee = repository.query(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10));
assertEquals(2, byAssignee.totalElements());
assertEquals(List.of(t3, t1), byAssignee.content());
Page<ApprovalTask> byStatus = repository.query(TaskQuery.any().withStatus(TaskStatus.PENDING), new PageRequest(0, 10));
assertEquals(List.of(t2, t1), byStatus.content());
Page<ApprovalTask> 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<ApprovalTask> 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<ApprovalTask> 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);
}
}
@@ -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<ProcessDefinition> 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<ProcessDefinition> 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")));
}
}
@@ -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<ProcessInstance> byDefinition = repository.query(InstanceQuery.any().withDefinitionId("leave"),
new PageRequest(0, 10));
assertEquals(2, byDefinition.totalElements());
assertEquals(List.of(i2, i1), byDefinition.content());
Page<ProcessInstance> byStatus = repository.query(InstanceQuery.any().withStatus(ProcessStatus.WITHDRAWN),
new PageRequest(0, 10));
assertEquals(List.of(i3), byStatus.content());
Page<ProcessInstance> 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<ProcessInstance> 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());
}
}
@@ -1,13 +1,18 @@
package com.jetlumen.ordo.storage.jdbc; package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalTask; 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.api.repository.ApprovalTaskRepository;
import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper;
import java.sql.Connection; import java.sql.Connection;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.ResultSet; import java.sql.ResultSet;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.Timestamp;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
@@ -36,6 +41,9 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
private static final String SELECT_PENDING_BY_INSTANCE = private static final String SELECT_PENDING_BY_INSTANCE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND instance_id = ?" "SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND instance_id = ?"
+ " ORDER BY created_at, 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; private final JdbcConnectionProvider connectionProvider;
@@ -111,6 +119,69 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
} }
} }
@Override
public Page<ApprovalTask> 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<ApprovalTask> 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<String> conditions = new ArrayList<>();
List<Object> 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<ApprovalTask> findOne(String sql, String parameter) { private Optional<ApprovalTask> findOne(String sql, String parameter) {
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) { try (PreparedStatement select = connection.prepareStatement(sql)) {
@@ -3,6 +3,8 @@ package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ApprovalStep; import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.StepTransition; 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.api.repository.ProcessDefinitionRepository;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalStepMapper.StepRow; 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 = private static final String SELECT_TRANSITIONS =
"SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition " "SELECT from_step_id, to_step_id, condition_key, priority FROM ordo_step_transition "
+ "WHERE definition_id = ? ORDER BY from_step_id, priority"; + "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; private final JdbcConnectionProvider connectionProvider;
@@ -112,43 +116,7 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
name = resultSet.getString("name"); name = resultSet.getString("name");
} }
} }
List<StepRow> stepRows = new ArrayList<>(); return Optional.of(assemble(connection, definitionId, name));
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<String, List<String>> 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<ApprovalStep> steps = new ArrayList<>();
for (StepRow row : stepRows) {
List<String> candidates = candidatesByStep.getOrDefault(row.stepId(), List.of());
steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(),
row.actionKey()));
}
List<StepTransition> 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));
} catch (SQLException e) { } catch (SQLException e) {
throw new JdbcStorageException("failed to load definition: " + definitionId, e); throw new JdbcStorageException("failed to load definition: " + definitionId, e);
} finally { } finally {
@@ -156,6 +124,74 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
} }
} }
@Override
public Page<ProcessDefinition> 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<ProcessDefinition> content = new ArrayList<>();
try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITIONS_PAGE)) {
select.setInt(1, pageRequest.size());
select.setInt(2, pageRequest.offset());
List<Map.Entry<String, String>> idsAndNames = new ArrayList<>();
try (ResultSet resultSet = select.executeQuery()) {
while (resultSet.next()) {
idsAndNames.add(Map.entry(resultSet.getString("id"), resultSet.getString("name")));
}
}
for (Map.Entry<String, String> 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<StepRow> 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<String, List<String>> 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<ApprovalStep> steps = new ArrayList<>();
for (StepRow row : stepRows) {
List<String> candidates = candidatesByStep.getOrDefault(row.stepId(), List.of());
steps.add(new ApprovalStep(row.stepId(), row.stepName(), candidates, row.policy(), row.kind(),
row.actionKey()));
}
List<StepTransition> 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 { private static boolean definitionExists(Connection connection, String definitionId) throws SQLException {
try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITION)) { try (PreparedStatement select = connection.prepareStatement(SELECT_DEFINITION)) {
select.setString(1, definitionId); select.setString(1, definitionId);
@@ -2,13 +2,20 @@ package com.jetlumen.ordo.storage.jdbc;
import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus; 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.api.repository.ProcessInstanceRepository;
import com.jetlumen.ordo.storage.jdbc.PageSupport.WhereClause;
import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper;
import java.sql.Connection; import java.sql.Connection;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.ResultSet; import java.sql.ResultSet;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.Timestamp;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.Optional; import java.util.Optional;
@@ -107,4 +114,62 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
connectionProvider.close(connection); connectionProvider.close(connection);
} }
} }
@Override
public Page<ProcessInstance> 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<ProcessInstance> 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<String> conditions = new ArrayList<>();
List<Object> 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);
}
} }
@@ -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<Object> params) {
}
static long count(Connection connection, String countSql, List<Object> 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<Object> params) throws SQLException {
int index = 1;
for (Object param : params) {
statement.setObject(index++, param);
}
return index;
}
}
@@ -8,6 +8,9 @@ import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus; import com.jetlumen.ordo.api.ProcessStatus;
import com.jetlumen.ordo.api.TaskAction; import com.jetlumen.ordo.api.TaskAction;
import com.jetlumen.ordo.api.TaskStatus; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -95,6 +98,41 @@ class JdbcApprovalTaskRepositoryTest {
assertEquals(List.of(henryTask), repository.findPendingByAssignee("henry")); 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<ApprovalTask> byAssignee = repository.query(TaskQuery.any().withAssignee("maria"), new PageRequest(0, 10));
assertEquals(2, byAssignee.totalElements());
assertEquals(List.of(t3, t1), byAssignee.content());
Page<ApprovalTask> byDefinition = repository.query(TaskQuery.any().withDefinitionId("leave"), new PageRequest(0, 10));
assertEquals(3, byDefinition.totalElements());
Page<ApprovalTask> 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<ApprovalTask> pageTwo = repository.query(TaskQuery.any(), new PageRequest(1, 2));
assertEquals(List.of(t2, t1), pageTwo.content());
assertFalse(pageTwo.hasNext());
}
@Test @Test
void onlyOneOfTwoConcurrentCompletionsWins() throws Exception { void onlyOneOfTwoConcurrentCompletionsWins() throws Exception {
ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT); ApprovalTask pending = pendingTask("task-1", "inst-1", "manager", "maria", CREATED_AT);
@@ -4,6 +4,8 @@ import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.StepKind; import com.jetlumen.ordo.api.StepKind;
import com.jetlumen.ordo.api.StepTransition; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -91,4 +93,21 @@ class JdbcProcessDefinitionRepositoryTest {
private static ProcessDefinition definition(String id, String name) { private static ProcessDefinition definition(String id, String name) {
return ProcessDefinition.linear(id, name, List.of(ApprovalStep.single("lead", "Lead approval", "lee"))); 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<ProcessDefinition> 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<ProcessDefinition> pageTwo = repository.findAll(new PageRequest(1, 2));
assertEquals(List.of("c-def"), pageTwo.content().stream().map(ProcessDefinition::id).toList());
assertFalse(pageTwo.hasNext());
}
} }
@@ -5,6 +5,9 @@ import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessInstance; import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -105,4 +108,39 @@ class JdbcProcessInstanceRepositoryTest {
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave", assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",
"alice", ProcessStatus.APPROVED, STARTED_AT, STARTED_AT, ProcessContext.empty()))); "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<ProcessInstance> byDefinition = repository.query(InstanceQuery.any().withDefinitionId("leave"),
new PageRequest(0, 10));
assertEquals(2, byDefinition.totalElements());
assertEquals(List.of(i2, i1), byDefinition.content());
Page<ProcessInstance> byInitiator = repository.query(InstanceQuery.any().withInitiator("alice"),
new PageRequest(0, 10));
assertEquals(List.of(i3, i1), byInitiator.content());
Page<ProcessInstance> 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<ProcessInstance> pageTwo = repository.query(InstanceQuery.any(), new PageRequest(1, 2));
assertEquals(List.of(i1), pageTwo.content());
assertFalse(pageTwo.hasNext());
}
} }