Compare commits

...
4 Commits
Author SHA1 Message Date
0264408andCursor e773badb23 feat: migrate Ordo schema with a dedicated Flyway history table
Stop rewriting spring.flyway.locations so host migrations stay independent and Boot 4 no longer needs starter-flyway for Ordo tables.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-22 10:47:43 +08:00
0264408andCursor 056ab7a806 docs: add JSON Schema for process definition documents
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-21 11:52:57 +08:00
0264408andCursor 02557a9b86 feat: schedule due processing from next dueAt instead of polling
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-20 16:54:13 +08:00
0264408andCursor 8d08a2840b docs: sync package layout and REST with current API
Host examples and README still described the pre-split packages and omitted optional REST.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-20 11:45:23 +08:00
23 changed files with 655 additions and 114 deletions
+18 -6
View File
@@ -1,6 +1,6 @@
# Ordo
轻量审批流程引擎。宿主通过 `OrdoEngine` 注册流程定义、发起实例、审批/驳回/转派/撤回/取消,并查询任务、实例与审计历史。引擎不绑定业务表单、用户体系或设计器 UI,当前也不自带 REST;业务数据放在 `ProcessContext` 里。
轻量审批流程引擎。宿主通过 `OrdoEngine` 注册流程定义、发起实例、审批/驳回/转派/撤回/取消,并查询任务、实例与审计历史。引擎不绑定业务表单、用户体系或设计器 UI。可选 REST(`ordo.rest.enabled`)。业务数据放在 `ProcessContext` 里。
要求 **Java 17+**。当前版本 `0.0.1-SNAPSHOT`。
@@ -14,7 +14,7 @@
| `ordo-spring-boot-starter` | Spring Boot 自动装配(JDBC + Flyway) |
| `ordo-example` | 内存引擎示例 |
存储实现通过 dialect 层支持 **PostgreSQL** 与 **MySQL**(测试可用 H2 PostgreSQL 兼容模式)。Starter **不携带** JDBC 驱动,宿主自行加入 `postgresql`、`mysql-connector-j` 或 `h2`。Spring Boot 4 还需额外引入 `spring-boot-starter-flyway`,否则迁移不会跑。
存储实现通过 dialect 层支持 **PostgreSQL** 与 **MySQL**(测试可用 H2 PostgreSQL 兼容模式)。Starter **不携带** JDBC 驱动,宿主自行加入 `postgresql`、`mysql-connector-j` 或 `h2`。表结构由 Ordo 自己的 Flyway 迁移(`ordo_schema_history`),与宿主 `spring.flyway` 无关。
## 能力
@@ -30,7 +30,7 @@
- 审计时间线:`ProcessEvent` + `queryHistory`
- 扩展点:`AssigneeResolver`、`RoutingCondition`、`ActionHandler`、`NamedCondition`、`NamedAction`、`OrdoEventListener`
开发计划:可选 REST + 目录 SPI。设计器为独立产品(不进本仓库),待 REST 与目录之后。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。
开发计划:独立设计器(不进本仓库)。多租户 **暂不在计划中**。见 [docs/roadmap.md](docs/roadmap.md)。
详细用法(定义 JSON、扩展点、异常、查询、ACTION/审计语义)见 **[docs/usage.md](docs/usage.md)**。对外行为变更时同步更新该文档。
@@ -45,6 +45,12 @@
```
```java
import com.jetlumen.ordo.api.OrdoEngine;
import com.jetlumen.ordo.api.definition.ApprovalStep;
import com.jetlumen.ordo.api.definition.ProcessDefinition;
import com.jetlumen.ordo.api.runtime.ProcessContext;
import com.jetlumen.ordo.core.InMemoryOrdoEngine;
OrdoEngine ordo = new InMemoryOrdoEngine();
ordo.publish(ProcessDefinition.linear("leave-request", "Leave request", List.of(
ApprovalStep.single("manager", "Manager approval", "maria"),
@@ -83,8 +89,7 @@ ordo.approve(manager.id(), "maria", "ok");
- `kind` 默认 `APPROVAL`;ACTION 用 `"action"` 作为 handler 查找键;`PARALLEL` 用 `branches`(结构化并行块,见 [docs/usage.md](docs/usage.md))。
- 带 `when` 的边先按 `priority` 匹配,都未命中再走无条件边;`to: null` 表示结束。
- 代码侧可用 `ProcessDefinitionParser.fromJson(...)`。
- `replace` 会整体替换同 id 定义;存在 `RUNNING` 实例时拒绝替换。
- 代码侧可用 `ProcessDefinitionParser.fromJson(...)`(`com.jetlumen.ordo.api.definition`)。
## Spring Boot
@@ -96,7 +101,7 @@ ordo.approve(manager.id(), "maria", "ok");
</dependency>
```
需要 `DataSource`。存在 `DataSource` 且 `ordo.enabled` 不为 `false` 时装配 JDBC 引擎。
需要 `DataSource`。存在 `DataSource` 且 `ordo.enabled` 不为 `false` 时装配 JDBC 引擎。REST 另需宿主的 Web starter,且 `ordo.rest.enabled=true`。
```yaml
ordo:
@@ -105,6 +110,12 @@ ordo:
dialect: # 可选 postgresql / mysql;空则按 DataSource 探测
definitions:
location: classpath*:ordo/*.json # 默认值;启动时 publish 加载
due:
poll-ms: 0 # >0 启用调度;值为最长空闲,按下次 dueAt 唤醒,满批续拉
batch-size: 100
rest:
enabled: false
base-path: /ordo
```
宿主提供 Bean 即可覆盖默认值:
@@ -115,6 +126,7 @@ ordo:
| `NamedCondition` / `NamedAction` | 可多个;无对应门面时按 key 分发 |
| `RoutingCondition` | 无具名 condition 时始终匹配;有则未知 `ref` 为 false |
| `ActionHandler` | 无具名 action 时空操作;有则未知 key 抛错 |
| `OrdoCatalog` | REST 打开时从 Named* 投影;`assignees` 空 |
| `OrdoEventListener` | 可注册多个,提交后按顺序调用 |
## 运行时约定
+1
View File
@@ -471,6 +471,7 @@ components:
label:
type: string
ProcessDefinitionDocument:
description: Closed document shape is docs/process-definition.schema.json. Graph rules are enforced by ProcessDefinitionParser.
type: object
required: [id, name, steps, transitions]
properties:
+297
View File
@@ -0,0 +1,297 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"$id": "https://jetlumen.com/ordo/schema/process-definition.json",
"title": "Ordo process definition document",
"$comment": "Document shape for ProcessDefinitionParser.fromJson. Graph connectivity, unique ids, PARALLEL membership, due.goto branch rules, and predicate depth/leaf limits are enforced only by the parser.",
"type": "object",
"additionalProperties": false,
"required": ["id", "name", "steps", "transitions"],
"properties": {
"id": { "$ref": "#/$defs/nonEmptyString" },
"name": { "$ref": "#/$defs/nonEmptyString" },
"version": {
"type": "integer",
"minimum": 0,
"$comment": "Ignored by fromJson; written by toJson after publish."
},
"startStep": { "$ref": "#/$defs/nonEmptyString" },
"steps": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/topStep" }
},
"transitions": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/transition" }
}
},
"$defs": {
"nonEmptyString": {
"type": "string",
"pattern": ".*\\S.*"
},
"literal": {
"type": ["string", "number", "boolean"]
},
"topStep": {
"oneOf": [
{ "$ref": "#/$defs/approvalStep" },
{ "$ref": "#/$defs/actionStep" },
{ "$ref": "#/$defs/parallelStep" }
]
},
"leafStep": {
"oneOf": [
{ "$ref": "#/$defs/approvalStep" },
{ "$ref": "#/$defs/actionStep" }
]
},
"approvalStep": {
"type": "object",
"additionalProperties": false,
"required": ["id", "name", "candidates"],
"properties": {
"id": { "$ref": "#/$defs/nonEmptyString" },
"name": { "$ref": "#/$defs/nonEmptyString" },
"kind": { "const": "APPROVAL" },
"candidates": {
"type": "array",
"minItems": 1,
"uniqueItems": true,
"items": { "$ref": "#/$defs/nonEmptyString" }
},
"policy": { "type": "string", "enum": ["ANY", "ALL"] },
"due": { "$ref": "#/$defs/due" }
}
},
"actionStep": {
"type": "object",
"additionalProperties": false,
"required": ["id", "name", "kind", "action"],
"properties": {
"id": { "$ref": "#/$defs/nonEmptyString" },
"name": { "$ref": "#/$defs/nonEmptyString" },
"kind": { "const": "ACTION" },
"action": { "$ref": "#/$defs/nonEmptyString" },
"policy": { "type": "string", "enum": ["ANY", "ALL"] }
}
},
"parallelStep": {
"type": "object",
"additionalProperties": false,
"required": ["id", "name", "kind", "branches"],
"properties": {
"id": { "$ref": "#/$defs/nonEmptyString" },
"name": { "$ref": "#/$defs/nonEmptyString" },
"kind": { "const": "PARALLEL" },
"policy": { "type": "string", "enum": ["ANY", "ALL"] },
"branches": {
"type": "array",
"minItems": 2,
"items": { "$ref": "#/$defs/branch" }
}
}
},
"branch": {
"type": "object",
"additionalProperties": false,
"required": ["id", "steps", "transitions"],
"properties": {
"id": { "$ref": "#/$defs/nonEmptyString" },
"steps": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/leafStep" }
},
"transitions": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/transition" }
}
}
},
"due": {
"oneOf": [
{
"type": "object",
"additionalProperties": false,
"required": ["after", "then", "to"],
"properties": {
"after": { "$ref": "#/$defs/isoDuration" },
"then": { "const": "reassign" },
"to": { "$ref": "#/$defs/nonEmptyString" }
}
},
{
"type": "object",
"additionalProperties": false,
"required": ["after", "then"],
"properties": {
"after": { "$ref": "#/$defs/isoDuration" },
"then": { "const": "notify" },
"action": { "$ref": "#/$defs/nonEmptyString" }
}
},
{
"type": "object",
"additionalProperties": false,
"required": ["after", "then", "to"],
"properties": {
"after": { "$ref": "#/$defs/isoDuration" },
"then": { "const": "goto" },
"to": { "$ref": "#/$defs/nonEmptyString" }
}
}
]
},
"isoDuration": {
"type": "string",
"minLength": 2,
"$comment": "java.time.Duration.parse; must be positive after parse."
},
"transition": {
"type": "object",
"additionalProperties": false,
"required": ["from"],
"properties": {
"from": { "$ref": "#/$defs/nonEmptyString" },
"to": {
"type": ["string", "null"],
"minLength": 1
},
"when": { "$ref": "#/$defs/when" },
"priority": { "type": "integer" }
}
},
"when": {
"oneOf": [
{ "$ref": "#/$defs/refWhen" },
{ "$ref": "#/$defs/predicate" }
]
},
"refWhen": {
"type": "object",
"additionalProperties": false,
"required": ["ref"],
"properties": {
"ref": { "$ref": "#/$defs/nonEmptyString" },
"args": {
"type": "object",
"additionalProperties": true
}
}
},
"predicate": {
"oneOf": [
{ "$ref": "#/$defs/compareEq" },
{ "$ref": "#/$defs/compareNe" },
{ "$ref": "#/$defs/compareGt" },
{ "$ref": "#/$defs/compareGte" },
{ "$ref": "#/$defs/compareLt" },
{ "$ref": "#/$defs/compareLte" },
{ "$ref": "#/$defs/inPredicate" },
{ "$ref": "#/$defs/andPredicate" },
{ "$ref": "#/$defs/orPredicate" },
{ "$ref": "#/$defs/notPredicate" }
]
},
"comparePair": {
"type": "array",
"minItems": 2,
"maxItems": 2,
"prefixItems": [
{ "$ref": "#/$defs/nonEmptyString" },
{ "$ref": "#/$defs/literal" }
]
},
"compareEq": {
"type": "object",
"additionalProperties": false,
"required": ["eq"],
"properties": { "eq": { "$ref": "#/$defs/comparePair" } }
},
"compareNe": {
"type": "object",
"additionalProperties": false,
"required": ["ne"],
"properties": { "ne": { "$ref": "#/$defs/comparePair" } }
},
"compareGt": {
"type": "object",
"additionalProperties": false,
"required": ["gt"],
"properties": { "gt": { "$ref": "#/$defs/comparePair" } }
},
"compareGte": {
"type": "object",
"additionalProperties": false,
"required": ["gte"],
"properties": { "gte": { "$ref": "#/$defs/comparePair" } }
},
"compareLt": {
"type": "object",
"additionalProperties": false,
"required": ["lt"],
"properties": { "lt": { "$ref": "#/$defs/comparePair" } }
},
"compareLte": {
"type": "object",
"additionalProperties": false,
"required": ["lte"],
"properties": { "lte": { "$ref": "#/$defs/comparePair" } }
},
"inPredicate": {
"type": "object",
"additionalProperties": false,
"required": ["in"],
"properties": {
"in": {
"type": "array",
"minItems": 2,
"maxItems": 2,
"prefixItems": [
{ "$ref": "#/$defs/nonEmptyString" },
{
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/literal" }
}
]
}
}
},
"andPredicate": {
"type": "object",
"additionalProperties": false,
"required": ["and"],
"properties": {
"and": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/predicate" }
}
}
},
"orPredicate": {
"type": "object",
"additionalProperties": false,
"required": ["or"],
"properties": {
"or": {
"type": "array",
"minItems": 1,
"items": { "$ref": "#/$defs/predicate" }
}
}
},
"notPredicate": {
"type": "object",
"additionalProperties": false,
"required": ["not"],
"properties": {
"not": { "$ref": "#/$defs/predicate" }
}
}
}
}
+1
View File
@@ -20,6 +20,7 @@
- 管理员/系统取消 `cancel` / `CANCELLED`
- 可选 REST(autoconfigure 条件装配)+ 目录 SPI `OrdoCatalog`
- `NamedAction` / `NamedCondition` 注册与官方 key 分发;默认 Catalog 从具名 Bean 投影
- API 分包:`api.definition` / `runtime` / `spi` / `util`
## 开发计划(确定要做)
+27 -6
View File
@@ -22,7 +22,19 @@ Ordo 是嵌入宿主进程的审批引擎,入口是 `OrdoEngine`。
| `ordo-spring-boot-starter` | Spring Boot 自动装配 |
| `ordo-example` | `LeaveRequestExample` 内存演示 |
Starter **不携带** JDBC 驱动。生产按库添加 `org.postgresql:postgresql` 或 `com.mysql:mysql-connector-j`。Spring Boot 4 还需 `spring-boot-starter-flyway`,否则 Flyway 迁移不会执行。
Java 包(模块未变):
| 包 | 内容 |
|---|---|
| `com.jetlumen.ordo.api` | `OrdoEngine`、`TransactionExecutor` |
| `com.jetlumen.ordo.api.definition` | 图:定义、步骤、边、`when`、PARALLEL |
| `com.jetlumen.ordo.api.runtime` | 实例、任务、事件、ACTION 执行、token、`ProcessRuntime` |
| `com.jetlumen.ordo.api.spi` | `ActionHandler`、`NamedAction`、`RoutingCondition`、`NamedCondition`、`AssigneeResolver`、`OrdoCatalog`、`OrdoEventListener` |
| `com.jetlumen.ordo.api.util` | `Texts`、`Jsons` |
| `com.jetlumen.ordo.api.exception` / `query` / `repository` | 异常、分页查询、存储端口 |
| `com.jetlumen.ordo.core.spi` | `DispatchingActionHandler`、`DispatchingRoutingCondition`、`RegistryOrdoCatalog` |
Starter **不携带** JDBC 驱动。生产按库添加 `org.postgresql:postgresql` 或 `com.mysql:mysql-connector-j`。Ordo 用独立 Flyway 建表,不依赖宿主 `spring.flyway`。
先 `mvn install` 本仓库,宿主再依赖 `0.0.1-SNAPSHOT`。
@@ -60,10 +72,13 @@ ordo:
enabled: true
jdbc:
dialect: # 可选 postgresql / mysql;空则按 DataSource 探测
migrate: true
history-table: ordo_schema_history
definitions:
location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 publish
due:
poll-ms: 0 # >0 时轮询 processDue;默认不调度
poll-ms: 0 # >0 启用调度;值为最长空闲,按下次 dueAt 唤醒,满批续拉
batch-size: 100
rest:
enabled: false
base-path: /ordo
@@ -85,7 +100,7 @@ ordo:
未提供 `NamedAction` 且未覆盖 `ActionHandler` 时,ACTION 步骤仍会推进流程,但 handler 什么都不做。不要同时提供门面 Bean 与对应 `Named*`(门面优先,具名 Bean 不参与运行时)。
内存引擎把 `DispatchingRoutingCondition.of(...)` / `DispatchingActionHandler.of(...)` 传入 `InMemoryOrdoEngine` 即可。
内存引擎把 `com.jetlumen.ordo.core.spi.DispatchingRoutingCondition.of(...)` / `DispatchingActionHandler.of(...)` 传入 `InMemoryOrdoEngine` 即可。
### 2.3 可选 REST
@@ -161,7 +176,7 @@ new ProcessDefinition("leave-request-routed", "Leave request",
### 3.2 JSON
`ProcessDefinitionParser.fromJson(String|InputStream)` / `toJson(ProcessDefinition)`。Spring 默认扫 `classpath*:ordo/*.json`。`toJson` 写出 `version` 与 `startStep`(当前步骤列表首位);`fromJson` 仍忽略 JSON 里的 `version`。
`ProcessDefinitionParser.fromJson(String|InputStream)` / `toJson(ProcessDefinition)`。Spring 默认扫 `classpath*:ordo/*.json`。`toJson` 写出 `version` 与 `startStep`(当前步骤列表首位);`fromJson` 仍忽略 JSON 里的 `version`。文档外形:[process-definition.schema.json](process-definition.schema.json);图连通、PARALLEL 约束、谓词深度/叶子上限仍以解析器为准。
```json
{
@@ -250,7 +265,8 @@ ordo.cancel(instance.id(), "admin", "政策变更");
- `approve` / `reject` / `reassign`:`actor` 必须等于该任务当前 `assignee`,否则 `UnauthorizedTaskOperationException`。
- `reassign`:仅 `PENDING` 任务;同一任务 id,办理人改为 `newAssignee`,不推进步骤。`newAssignee` 不可空白、不可等于当前 `assignee`,且同一步不能已有该人的 `PENDING` 任务,否则 `IllegalArgumentException`。不经过 `AssigneeResolver`。人工转派不改 `dueAt`。
- `processDue(limit)`:认领 `dueAt <= now` 的 PENDING 任务(`limit > 0`),按步上 `due.then` 执行:`reassign` 换办理人(`to` 走 `AssigneeResolver`)、`notify` 可选 `ActionHandler`、`goto` 跳过当前步 PENDING 并进入 `to` 步骤(`to` 必须是步骤 id)。每种策略对一张任务最多成功一次(清空 `dueAt`)。引擎无后台线程;Spring 下 `ordo.due.poll-ms > 0` 才轮询。
- `processDue(limit)`:认领 `dueAt <= now` 的 PENDING 任务(`limit > 0`),按步上 `due.then` 执行:`reassign` 换办理人(`to` 走 `AssigneeResolver`)、`notify` 可选 `ActionHandler`、`goto` 跳过当前步 PENDING 并进入 `to` 步骤(`to` 必须是步骤 id)。每种策略对一张任务最多成功一次(清空 `dueAt`)。引擎无后台线程;Spring 下 `ordo.due.poll-ms > 0` 才调度:按下次 `dueAt` 唤醒,满批续拉,`poll-ms` 为最长空闲。
- `nextDueAt()`:PENDING 且仍有 `dueAt` 的最早到期时刻;没有则 empty。
- 任务非 `PENDING`:`TaskAlreadyCompletedException`。
- `withdraw`:仅 `initiator`,否则 `UnauthorizedInstanceOperationException`;实例非 `RUNNING`:`InstanceAlreadyCompletedException`。
- `cancel`:`actor` 非空即可,**不校验**是否发起人;实例须为 `RUNNING`,否则 `InstanceAlreadyCompletedException`。谁能调用由宿主决定。
@@ -314,6 +330,10 @@ v1:至少 2 条分支;禁止套娃 PARALLEL;join 固定 ALL;任一分支
**失败不回滚已提交的审批,不阻塞后续步骤,引擎不做重试。** 宿主用 `queryActionExecutions` 或 listener 自行补发。分发器遇到未知 `actionKey` 会抛 `IllegalArgumentException`,记为该次 ACTION `FAILED`。
```java
import com.jetlumen.ordo.api.runtime.ProcessRuntime;
import com.jetlumen.ordo.api.spi.NamedAction;
import org.springframework.stereotype.Component;
@Component
public class LeaveApprovedMail implements NamedAction {
@Override
@@ -376,6 +396,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| `approve` / `reject` | 办理当前 PENDING 任务 |
| `reassign` | 当前办理人把 PENDING 任务转给他人 |
| `processDue` | 认领并处理已到期 PENDING 任务 |
| `nextDueAt` | 下一笔 PENDING 任务的 `dueAt` |
| `withdraw` | 发起人撤回 |
| `cancel` | 管理员/系统取消(引擎不鉴权角色) |
| `find*` | 按 id / 待办索引读取 |
@@ -401,7 +422,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
## 12. 存储
Flyway 脚本按方言分目录:`db/postgresql/migration`、`db/mysql/migration`。未设置 `spring.flyway.locations` 时,starter 按探测到的方言指向对应目录。已有 `flyway_schema_history` 的开发库按版本追加迁移(PARALLEL 为 `V3__parallel_blocks.sql`)。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤(含 `parent_step_id` / `branch_id`)/候选人/转移、实例(含 `definition_version`)、任务、并行令牌 `ordo_instance_token`、`ordo_process_event`、`ordo_action_execution`。
Flyway 脚本按方言分目录:`db/postgresql/migration`、`db/mysql/migration`。Starter 在启动时用独立 Flyway 执行这些脚本,历史表默认 `ordo_schema_history`,不修改 `spring.flyway.locations`。若宿主也启用了 Spring Boot Flyway,Ordo 会在宿主 `flywayInitializer` 之后再 migrate,避免非空 schema 导致宿主失败。已有 `ordo_*` 表但尚无该历史表时,会 baseline 到当前脚本最高版本后再 migrate。`ordo.jdbc.migrate=false` 时不执行。表包括流程头 `ordo_process`、按 `(id, version)` 存储的定义/步骤(含 `parent_step_id` / `branch_id`)/候选人/转移、实例(含 `definition_version`)、任务、并行令牌 `ordo_instance_token`、`ordo_process_event`、`ordo_action_execution`。
自定义方言:实现 `SqlDialect` 并用 `META-INF/services` 注册,或提供 `SqlDialect` Bean。新增列时每个已支持方言目录各加一条迁移。
@@ -11,6 +11,7 @@ import com.jetlumen.ordo.api.runtime.ProcessContext;
import com.jetlumen.ordo.api.runtime.ProcessEvent;
import com.jetlumen.ordo.api.runtime.ProcessInstance;
import java.time.Instant;
import java.util.List;
import java.util.Optional;
@@ -40,6 +41,9 @@ public interface OrdoEngine {
/** Claims and processes up to {@code limit} overdue pending tasks. {@code limit} must be positive. */
int processDue(int limit);
/** Earliest {@code dueAt} among pending tasks, or empty if none are scheduled. */
Optional<Instant> nextDueAt();
default ProcessInstance withdraw(String instanceId, String actor) {
return withdraw(instanceId, actor, null);
}
@@ -43,6 +43,9 @@ public interface ApprovalTaskRepository {
List<ApprovalTask> findDuePending(Instant now, int limit);
/** Earliest {@code dueAt} among pending tasks, or empty if none are scheduled. */
Optional<Instant> findNextDueAt();
boolean claimIfDue(String taskId, String expectedAssignee, Instant now);
/** Paginated, filterable query; results are ordered newest-first (created_at desc). */
@@ -345,6 +345,11 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return processed;
}
@Override
public synchronized Optional<Instant> nextDueAt() {
return taskRepository.findNextDueAt();
}
private boolean escalateDueTask(ApprovalTask overdue, Instant now, List<PendingAction> queued,
List<ProcessEvent> events) {
if (!taskRepository.claimIfDue(overdue.id(), overdue.assignee(), now)) {
@@ -23,6 +23,7 @@ import com.jetlumen.ordo.core.repository.InMemoryProcessHistoryRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository;
import java.time.Clock;
import java.time.Instant;
import java.util.List;
import java.util.Optional;
@@ -181,4 +182,9 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
public int processDue(int limit) {
return delegate.processDue(limit);
}
@Override
public Optional<Instant> nextDueAt() {
return delegate.nextDueAt();
}
}
@@ -113,6 +113,15 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
.toList();
}
@Override
public synchronized Optional<Instant> findNextDueAt() {
return tasks.values().stream()
.filter(task -> task.status() == TaskStatus.PENDING)
.map(ApprovalTask::dueAt)
.filter(Objects::nonNull)
.min(Comparator.naturalOrder());
}
@Override
public synchronized boolean claimIfDue(String taskId, String expectedAssignee, Instant now) {
ApprovalTask current = tasks.get(taskId);
@@ -5,24 +5,35 @@ import org.springframework.context.SmartLifecycle;
import java.lang.System.Logger;
import java.lang.System.Logger.Level;
import java.time.Clock;
import java.time.Instant;
import java.util.Objects;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
/** Optional poller that calls {@link OrdoEngine#processDue(int)} when {@code ordo.due.poll-ms} is positive. */
/** Optional scheduler that drains {@link OrdoEngine#processDue(int)} when {@code ordo.due.poll-ms} is positive. */
public final class OrdoDuePoller implements SmartLifecycle {
private static final Logger LOG = System.getLogger("ordo");
private static final int BATCH_SIZE = 100;
private final OrdoEngine ordoEngine;
private final Clock clock;
private final long pollMs;
private final int batchSize;
private ScheduledExecutorService executor;
private ScheduledFuture<?> scheduled;
private volatile boolean running;
private volatile boolean ticking;
public OrdoDuePoller(OrdoEngine ordoEngine, long pollMs) {
public OrdoDuePoller(OrdoEngine ordoEngine, Clock clock, long pollMs, int batchSize) {
this.ordoEngine = Objects.requireNonNull(ordoEngine, "ordoEngine must not be null");
this.clock = Objects.requireNonNull(clock, "clock must not be null");
this.pollMs = pollMs;
if (batchSize <= 0) {
throw new IllegalArgumentException("batch size must be positive");
}
this.batchSize = batchSize;
}
@Override
@@ -40,16 +51,49 @@ public final class OrdoDuePoller implements SmartLifecycle {
thread.setDaemon(true);
return thread;
});
executor.scheduleWithFixedDelay(this::tick, pollMs, pollMs, TimeUnit.MILLISECONDS);
running = true;
schedule(0);
}
void wake() {
if (!running || ticking) {
return;
}
schedule(0);
}
private void tick() {
ticking = true;
try {
ordoEngine.processDue(BATCH_SIZE);
int processed;
do {
processed = ordoEngine.processDue(batchSize);
} while (processed == batchSize);
} catch (RuntimeException e) {
LOG.log(Level.WARNING, "processDue failed", e);
} finally {
if (running) {
schedule(nextDelayMs());
}
ticking = false;
}
}
private long nextDelayMs() {
Instant now = clock.instant();
return ordoEngine.nextDueAt()
.map(dueAt -> Math.min(Math.max(0L, dueAt.toEpochMilli() - now.toEpochMilli()), pollMs))
.orElse(pollMs);
}
private synchronized void schedule(long delayMs) {
if (!running || executor == null) {
return;
}
if (scheduled != null) {
scheduled.cancel(false);
}
scheduled = executor.schedule(this::tick, delayMs, TimeUnit.MILLISECONDS);
}
@Override
@@ -59,6 +103,7 @@ public final class OrdoDuePoller implements SmartLifecycle {
executor.shutdownNow();
executor = null;
}
scheduled = null;
}
@Override
@@ -0,0 +1,25 @@
package com.jetlumen.ordo.spring;
import com.jetlumen.ordo.api.runtime.ProcessEvent;
import com.jetlumen.ordo.api.runtime.ProcessEventType;
import com.jetlumen.ordo.api.spi.OrdoEventListener;
/** Forwards {@code TASK_CREATED} to {@link OrdoDuePoller} after the engine listener snapshot is taken. */
public final class OrdoDueWakeBridge implements OrdoEventListener {
private volatile OrdoDuePoller poller;
public void attach(OrdoDuePoller poller) {
this.poller = poller;
}
@Override
public void onEvent(ProcessEvent event) {
if (event.type() != ProcessEventType.TASK_CREATED) {
return;
}
OrdoDuePoller current = poller;
if (current != null) {
current.wake();
}
}
}
@@ -51,9 +51,12 @@ public class OrdoProperties {
}
public static class Due {
/** Poll interval in milliseconds. {@code 0} disables scheduling. */
/** Enables scheduling when positive; also the maximum idle sleep in milliseconds. */
private long pollMs;
/** Tasks claimed per {@code processDue} call while draining. */
private int batchSize = 100;
public long getPollMs() {
return pollMs;
}
@@ -61,12 +64,26 @@ public class OrdoProperties {
public void setPollMs(long pollMs) {
this.pollMs = pollMs;
}
public int getBatchSize() {
return batchSize;
}
public void setBatchSize(int batchSize) {
this.batchSize = batchSize;
}
}
public static class Jdbc {
/** Explicit dialect id ({@code postgresql}, {@code mysql}). Empty means detect from the DataSource. */
private String dialect;
/** When true, Ordo runs its own Flyway against dialect locations. */
private boolean migrate = true;
/** Flyway history table used only by Ordo (not {@code flyway_schema_history}). */
private String historyTable = "ordo_schema_history";
public String getDialect() {
return dialect;
}
@@ -74,6 +91,23 @@ public class OrdoProperties {
public void setDialect(String dialect) {
this.dialect = dialect;
}
public boolean isMigrate() {
return migrate;
}
public void setMigrate(boolean migrate) {
this.migrate = migrate;
}
public String getHistoryTable() {
return historyTable;
}
public void setHistoryTable(String historyTable) {
this.historyTable = historyTable == null || historyTable.isBlank()
? "ordo_schema_history" : historyTable.trim();
}
}
public static class Rest {
@@ -0,0 +1,14 @@
package com.jetlumen.ordo.spring.jdbc;
import org.springframework.boot.sql.init.dependency.AbstractBeansOfTypeDatabaseInitializerDetector;
import java.util.Set;
/** Treats {@link OrdoSchemaMigrator} as database initialization for {@code @DependsOnDatabaseInitialization}. */
public class OrdoDatabaseInitializerDetector extends AbstractBeansOfTypeDatabaseInitializerDetector {
@Override
protected Set<Class<?>> getDatabaseInitializerBeanTypes() {
return Set.of(OrdoSchemaMigrator.class);
}
}
@@ -1,34 +1,22 @@
package com.jetlumen.ordo.spring.jdbc;
import com.jetlumen.ordo.spring.OrdoProperties;
import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialect;
import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects;
import org.flywaydb.core.Flyway;
import org.flywaydb.core.api.configuration.FluentConfiguration;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
import org.springframework.boot.autoconfigure.AutoConfigureBefore;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.core.env.Environment;
import javax.sql.DataSource;
import java.lang.reflect.Proxy;
/**
* Points Flyway at the dialect migration directory before Flyway runs.
* Host {@code spring.flyway.locations} is left unchanged when set.
*
* <p>Does not inject {@link DataSource} or {@link SqlDialect} at bean-creation time
* (that would cycle with DataSource → Flyway → customizer). Dialect is resolved inside
* {@code customize} from the FluentConfiguration DataSource / {@code ordo.jdbc.dialect}.
*
* <p>FlywayConfigurationCustomizer moved between Boot 3 and Boot 4; a reflective
* {@link FactoryBean} supplies a proxy for whichever type is on the classpath.
* Migrates Ordo tables with a dedicated Flyway history table, independent of
* {@code spring.flyway}. Runs after the host {@code flywayInitializer} when present so a
* non-empty schema does not break host migrate.
*/
@AutoConfiguration
@ConditionalOnProperty(prefix = "ordo", name = "enabled", havingValue = "true", matchIfMissing = true)
@@ -36,89 +24,19 @@ import java.lang.reflect.Proxy;
@ConditionalOnBean(DataSource.class)
@AutoConfigureAfter(name = {
"org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration",
"org.springframework.boot.jdbc.autoconfigure.DataSourceAutoConfiguration"
})
@AutoConfigureBefore(name = {
"org.springframework.boot.jdbc.autoconfigure.DataSourceAutoConfiguration",
"org.springframework.boot.autoconfigure.flyway.FlywayAutoConfiguration",
"org.springframework.boot.flyway.autoconfigure.FlywayAutoConfiguration"
})
@EnableConfigurationProperties(OrdoProperties.class)
public class OrdoFlywayAutoConfiguration {
private static final String[] FLYWAY_CUSTOMIZER_TYPES = {
"org.springframework.boot.flyway.autoconfigure.FlywayConfigurationCustomizer",
"org.springframework.boot.autoconfigure.flyway.FlywayConfigurationCustomizer"
};
@Bean
@ConditionalOnClass(name = "org.flywaydb.core.api.configuration.FluentConfiguration")
public FactoryBean<Object> ordoFlywayConfigurationCustomizer(
OrdoProperties ordoProperties, Environment environment) {
Class<?> customizerType = resolveFlywayCustomizerType();
if (customizerType == null) {
return null;
}
return new FlywayLocationsCustomizerFactoryBean(customizerType, ordoProperties, environment);
}
static Class<?> resolveFlywayCustomizerType() {
ClassLoader classLoader = OrdoFlywayAutoConfiguration.class.getClassLoader();
for (String name : FLYWAY_CUSTOMIZER_TYPES) {
try {
return Class.forName(name, false, classLoader);
} catch (ClassNotFoundException ignored) {
// Boot 3 vs Boot 4
}
}
return null;
}
private static final class FlywayLocationsCustomizerFactoryBean implements FactoryBean<Object> {
private final Class<?> customizerType;
private final OrdoProperties ordoProperties;
private final Environment environment;
private FlywayLocationsCustomizerFactoryBean(Class<?> customizerType, OrdoProperties ordoProperties,
Environment environment) {
this.customizerType = customizerType;
this.ordoProperties = ordoProperties;
this.environment = environment;
}
@Override
public Object getObject() {
return Proxy.newProxyInstance(customizerType.getClassLoader(), new Class<?>[] {customizerType},
(proxy, method, args) -> {
String name = method.getName();
if ("customize".equals(name) && args != null && args.length == 1) {
applyLocations((FluentConfiguration) args[0]);
return null;
}
if ("equals".equals(name)) {
return proxy == args[0];
}
if ("hashCode".equals(name)) {
return System.identityHashCode(proxy);
}
if ("toString".equals(name)) {
return "OrdoFlywayLocationsCustomizer";
}
throw new UnsupportedOperationException(method.toString());
});
}
private void applyLocations(FluentConfiguration configuration) {
if (environment.containsProperty("spring.flyway.locations")) {
return;
}
SqlDialect dialect = SqlDialects.resolve(configuration.getDataSource(),
ordoProperties.getJdbc().getDialect());
configuration.locations(dialect.flywayLocations());
}
@Override
public Class<?> getObjectType() {
return customizerType;
public OrdoSchemaMigrator ordoSchemaMigrator(DataSource dataSource, OrdoProperties ordoProperties,
BeanFactory beanFactory) {
if (beanFactory.containsBean("flywayInitializer")) {
beanFactory.getBean("flywayInitializer");
}
return new OrdoSchemaMigrator(dataSource, ordoProperties);
}
}
@@ -19,6 +19,7 @@ import com.jetlumen.ordo.core.spi.DispatchingActionHandler;
import com.jetlumen.ordo.core.spi.DispatchingRoutingCondition;
import com.jetlumen.ordo.spring.OrdoDefinitionLoader;
import com.jetlumen.ordo.spring.OrdoDuePoller;
import com.jetlumen.ordo.spring.OrdoDueWakeBridge;
import com.jetlumen.ordo.spring.OrdoProperties;
import com.jetlumen.ordo.storage.jdbc.JdbcActionExecutionRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcApprovalTaskRepository;
@@ -152,6 +153,12 @@ public class OrdoJdbcAutoConfiguration {
return new JdbcInstanceTokenRepository(connectionProvider, ordoSqlDialect);
}
@Bean
@ConditionalOnMissingBean
public OrdoDueWakeBridge ordoDueWakeBridge() {
return new OrdoDueWakeBridge();
}
@Bean
@ConditionalOnMissingBean
public OrdoEngine ordoEngine(Clock ordoClock,
@@ -181,7 +188,11 @@ public class OrdoJdbcAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public OrdoDuePoller ordoDuePoller(OrdoEngine ordoEngine, OrdoProperties ordoProperties) {
return new OrdoDuePoller(ordoEngine, ordoProperties.getDue().getPollMs());
public OrdoDuePoller ordoDuePoller(OrdoEngine ordoEngine, Clock ordoClock, OrdoProperties ordoProperties,
OrdoDueWakeBridge ordoDueWakeBridge) {
OrdoDuePoller poller = new OrdoDuePoller(ordoEngine, ordoClock, ordoProperties.getDue().getPollMs(),
ordoProperties.getDue().getBatchSize());
ordoDueWakeBridge.attach(poller);
return poller;
}
}
@@ -0,0 +1,86 @@
package com.jetlumen.ordo.spring.jdbc;
import com.jetlumen.ordo.spring.OrdoProperties;
import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects;
import org.flywaydb.core.Flyway;
import org.flywaydb.core.api.MigrationInfo;
import org.flywaydb.core.api.MigrationVersion;
import org.flywaydb.core.api.configuration.FluentConfiguration;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.Locale;
/**
* Runs Ordo Flyway scripts independently of the host {@code spring.flyway} instance.
* Presence of this bean marks Ordo schema initialization for
* {@code @DependsOnDatabaseInitialization}.
*/
public final class OrdoSchemaMigrator {
private static final String PROCESS_TABLE = "ordo_process";
public OrdoSchemaMigrator(DataSource dataSource, OrdoProperties properties) {
OrdoProperties.Jdbc jdbc = properties.getJdbc();
if (!jdbc.isMigrate()) {
return;
}
String[] locations = SqlDialects.resolve(dataSource, jdbc.getDialect()).flywayLocations();
String historyTable = jdbc.getHistoryTable();
FluentConfiguration configuration = new FluentConfiguration(OrdoSchemaMigrator.class.getClassLoader())
.dataSource(dataSource)
.locations(locations)
.table(historyTable);
applyBaselineIfNeeded(configuration, dataSource, historyTable);
configuration.load().migrate();
}
private static void applyBaselineIfNeeded(FluentConfiguration configuration, DataSource dataSource,
String historyTable) {
try (Connection connection = dataSource.getConnection()) {
DatabaseMetaData metaData = connection.getMetaData();
if (tableExists(metaData, connection, historyTable)) {
return;
}
// Host tables / flyway_schema_history make the schema non-empty; Flyway then requires
// baseline before migrate. Existing Ordo tables baseline to latest (skip scripts);
// otherwise baseline at 0 so V1+ still run.
if (tableExists(metaData, connection, PROCESS_TABLE)) {
MigrationVersion latest = latestVersion(configuration.load());
configuration.baselineOnMigrate(true).baselineVersion(latest);
} else {
configuration.baselineOnMigrate(true).baselineVersion(MigrationVersion.fromVersion("0"));
}
} catch (SQLException e) {
throw new IllegalStateException("failed to inspect schema before Ordo Flyway migrate", e);
}
}
private static boolean tableExists(DatabaseMetaData metaData, Connection connection, String table)
throws SQLException {
String catalog = connection.getCatalog();
String schema = connection.getSchema();
String[] names = {table, table.toLowerCase(Locale.ROOT), table.toUpperCase(Locale.ROOT)};
for (String name : names) {
try (ResultSet tables = metaData.getTables(catalog, schema, name, new String[] {"TABLE"})) {
if (tables.next()) {
return true;
}
}
}
return false;
}
private static MigrationVersion latestVersion(Flyway flyway) {
MigrationVersion latest = MigrationVersion.fromVersion("0");
for (MigrationInfo info : flyway.info().all()) {
if (info.getVersion() != null && info.getVersion().compareTo(latest) > 0) {
latest = info.getVersion();
}
}
return latest;
}
}
@@ -0,0 +1 @@
com.jetlumen.ordo.spring.jdbc.OrdoDatabaseInitializerDetector
@@ -73,7 +73,7 @@ class OrdoJdbcAutoConfigurationTest {
@Test
void honoursExplicitMysqlDialect() {
withDataSourceRunner.withPropertyValues("ordo.jdbc.dialect=mysql", "spring.flyway.enabled=false")
withDataSourceRunner.withPropertyValues("ordo.jdbc.dialect=mysql", "ordo.jdbc.migrate=false")
.run(context -> assertThat(context.getBean(SqlDialect.class).id()).isEqualTo("mysql"));
}
@@ -11,6 +11,7 @@ import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper;
import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
@@ -18,6 +19,7 @@ import java.sql.Timestamp;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
import java.util.Objects;
import java.util.Optional;
@@ -52,6 +54,8 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
private static final String SELECT_DUE_PENDING_BASE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL"
+ " AND due_at <= ? ORDER BY due_at, id";
private static final String SELECT_NEXT_DUE_AT =
"SELECT MIN(due_at) FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL";
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, t.due_at";
@@ -67,7 +71,7 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
this.dialect = Objects.requireNonNull(dialect, "dialect must not be null");
this.selectDuePending = dialect.limit(SELECT_DUE_PENDING_BASE, false);
this.selectDuePending = duePendingSql(this.connectionProvider, this.dialect);
}
@Override
@@ -180,6 +184,23 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
}
}
@Override
public Optional<Instant> findNextDueAt() {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(SELECT_NEXT_DUE_AT);
ResultSet resultSet = select.executeQuery()) {
if (!resultSet.next()) {
return Optional.empty();
}
Timestamp timestamp = resultSet.getTimestamp(1);
return timestamp == null ? Optional.empty() : Optional.of(timestamp.toInstant());
} catch (SQLException e) {
throw new JdbcStorageException("failed to query next due at", e);
} finally {
connectionProvider.close(connection);
}
}
@Override
public boolean claimIfDue(String taskId, String expectedAssignee, Instant now) {
Objects.requireNonNull(taskId, "taskId must not be null");
@@ -275,6 +296,24 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
}
}
private static String duePendingSql(JdbcConnectionProvider connectionProvider, SqlDialect dialect) {
String limited = dialect.limit(SELECT_DUE_PENDING_BASE, false);
return supportsSkipLocked(connectionProvider) ? dialect.forUpdateSkipLocked(limited) : limited;
}
private static boolean supportsSkipLocked(JdbcConnectionProvider connectionProvider) {
Connection connection = connectionProvider.getConnection();
try {
DatabaseMetaData metaData = connection.getMetaData();
String product = metaData.getDatabaseProductName();
return product != null && !product.toLowerCase(Locale.ROOT).contains("h2");
} catch (SQLException e) {
return false;
} finally {
connectionProvider.close(connection);
}
}
private List<ApprovalTask> findAll(String sql, String parameter) {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) {
@@ -16,4 +16,9 @@ abstract class LimitOffsetSqlDialect implements SqlDialect {
public final String forUpdate(String sql) {
return sql + " FOR UPDATE";
}
@Override
public final String forUpdateSkipLocked(String sql) {
return sql + " FOR UPDATE SKIP LOCKED";
}
}
@@ -25,5 +25,8 @@ public interface SqlDialect {
String forUpdate(String sql);
/** Appends {@code FOR UPDATE SKIP LOCKED}. */
String forUpdateSkipLocked(String sql);
String[] flywayLocations();
}
@@ -54,6 +54,7 @@ class SqlDialectsTest {
assertEquals("SELECT 1 LIMIT ? OFFSET ?", dialect.limit("SELECT 1"));
assertEquals("SELECT 1 LIMIT ?", dialect.limit("SELECT 1", false));
assertEquals("SELECT 1 FOR UPDATE", dialect.forUpdate("SELECT 1"));
assertEquals("SELECT 1 FOR UPDATE SKIP LOCKED", dialect.forUpdateSkipLocked("SELECT 1"));
assertEquals("classpath:db/postgresql/migration", dialect.flywayLocations()[0]);
}