feat: pass ProcessRuntime to host SPIs
Keep initiator and instance identity off ProcessContext variables so AssigneeResolver, RoutingCondition, and ActionHandler can read them from a dedicated runtime view. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+4
-4
@@ -270,7 +270,7 @@ ordo.cancel(instance.id(), "admin", "政策变更");
|
||||
| 通过 | 任一人 `approve` 即过步,其余 PENDING 变 `SKIPPED` | 全部 `approve` 才过步 |
|
||||
| 驳回 | 所有候选人都 `reject` 才否决实例 | 任一人 `reject` 即否决实例,其余 PENDING 变 `SKIPPED` |
|
||||
|
||||
候选人创建任务前会经过 `AssigneeResolver.resolve(candidate, step, context)`,例如把角色名解析成用户 id。默认实现原样返回 candidate。
|
||||
候选人创建任务前会经过 `AssigneeResolver.resolve(candidate, step, runtime)`,例如把角色名解析成用户 id。默认实现原样返回 candidate。`ProcessRuntime` 含 `instanceId`、`definitionId`、`definitionVersion`、`initiator`、当前 `stepId` 和业务 `context`。发起人等引擎元数据不写入 `ProcessContext.variables`。
|
||||
|
||||
连续 ACTION 步会在同一次提交后依次执行,上限 32 跳(**每个 PARALLEL 分支各自计数**),超出抛 `IllegalStateException`。
|
||||
|
||||
@@ -294,7 +294,7 @@ v1:至少 2 条分支;禁止套娃 PARALLEL;join 固定 ALL;任一分支
|
||||
|
||||
无条件边是默认分支;`priority` 只在同类边之间比较(多条条件边之间,或多条无条件边之间)。
|
||||
|
||||
`RoutingCondition` 只看到 `ProcessContext` 与 `args`,看不到任务意见,也没有流程定义 id。上下文在 `start` 时写入,运行中引擎不会改 context。未知 `ref` 由宿主返回 `false`,该边不匹配。
|
||||
`RoutingCondition.matches` 看到 `args` 与 `ProcessRuntime`(含业务 `context`、发起人、定义与当前步)。看不到任务意见。上下文在 `start` 时写入,运行中引擎不会改 context。未知 `ref` 由宿主返回 `false`,该边不匹配。内置谓词仍只读 `ProcessContext` 变量。
|
||||
|
||||
引擎只注入**一个** `RoutingCondition`。Spring 下多个该类型 Bean 会冲突。宿主用一个门面按 `ref` 分发;不要指望引擎按定义拆 bean。
|
||||
|
||||
@@ -304,7 +304,7 @@ v1:至少 2 条分支;禁止套娃 PARALLEL;join 固定 ALL;任一分支
|
||||
|
||||
1. 事务内插入 `ActionExecution`,状态 `PENDING`。
|
||||
2. 按转移进入下一步或结束实例(仍在同一事务)。
|
||||
3. 事务提交后调用 `ActionHandler.execute(actionKey, context)`。
|
||||
3. 事务提交后调用 `ActionHandler.execute(actionKey, runtime)`。
|
||||
4. 成功 → `SUCCESS` + 事件 `ACTION_SUCCEEDED`;失败 → `FAILED`(`errorMessage`)+ `ACTION_FAILED` + 日志 WARNING。
|
||||
|
||||
**失败不回滚已提交的审批,不阻塞后续步骤,引擎不做重试。** 宿主用 `queryActionExecutions` 或 listener 自行补发。
|
||||
@@ -313,7 +313,7 @@ v1:至少 2 条分支;禁止套娃 PARALLEL;join 固定 ALL;任一分支
|
||||
@Component
|
||||
public class MailActions implements ActionHandler {
|
||||
@Override
|
||||
public void execute(String actionKey, ProcessContext context) {
|
||||
public void execute(String actionKey, ProcessRuntime runtime) {
|
||||
switch (actionKey) {
|
||||
case "leave-submitted-mail" -> { /* ... */ }
|
||||
case "leave-approved-mail" -> { /* ... */ }
|
||||
|
||||
@@ -1,16 +1,16 @@
|
||||
package com.jetlumen.ordo.api;
|
||||
|
||||
/**
|
||||
* Executes a named action step against the running instance's context. Hosts supply a
|
||||
* Executes a named action step against the running instance. Hosts supply a
|
||||
* singleton implementation (same pattern as {@link RoutingCondition}); the database only stores
|
||||
* the {@code actionKey} string.
|
||||
*/
|
||||
@FunctionalInterface
|
||||
public interface ActionHandler {
|
||||
void execute(String actionKey, ProcessContext context);
|
||||
void execute(String actionKey, ProcessRuntime runtime);
|
||||
|
||||
static ActionHandler noop() {
|
||||
return (key, context) -> {
|
||||
return (key, runtime) -> {
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8,9 +8,9 @@ package com.jetlumen.ordo.api;
|
||||
*/
|
||||
@FunctionalInterface
|
||||
public interface AssigneeResolver {
|
||||
String resolve(String candidate, ApprovalStep step, ProcessContext context);
|
||||
String resolve(String candidate, ApprovalStep step, ProcessRuntime runtime);
|
||||
|
||||
static AssigneeResolver direct() {
|
||||
return (candidate, step, context) -> candidate;
|
||||
return (candidate, step, runtime) -> candidate;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
package com.jetlumen.ordo.api;
|
||||
|
||||
import java.util.Objects;
|
||||
|
||||
/**
|
||||
* Read-only engine metadata plus business {@link ProcessContext} for host SPIs.
|
||||
* Initiator and instance identity stay here; they are not copied into context variables.
|
||||
*/
|
||||
public record ProcessRuntime(String instanceId, String definitionId, int definitionVersion, String initiator,
|
||||
String stepId, ProcessContext context) {
|
||||
public ProcessRuntime {
|
||||
Texts.requireText(instanceId, "instance id");
|
||||
Texts.requireText(definitionId, "definition id");
|
||||
if (definitionVersion < 1) {
|
||||
throw new IllegalArgumentException("definition version must be at least 1");
|
||||
}
|
||||
Texts.requireText(initiator, "initiator");
|
||||
Texts.requireText(stepId, "step id");
|
||||
Objects.requireNonNull(context, "context must not be null");
|
||||
}
|
||||
|
||||
public static ProcessRuntime of(ProcessInstance instance, String stepId) {
|
||||
Objects.requireNonNull(instance, "instance must not be null");
|
||||
return new ProcessRuntime(instance.id(), instance.definitionId(), instance.definitionVersion(),
|
||||
instance.initiator(), stepId, instance.context());
|
||||
}
|
||||
}
|
||||
@@ -3,15 +3,15 @@ package com.jetlumen.ordo.api;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* Evaluates a named routing condition ({@code when.ref}) against the running instance's context.
|
||||
* Evaluates a named routing condition ({@code when.ref}) against the running instance.
|
||||
* Hosts supply a singleton implementation (same pattern as {@link AssigneeResolver}); the
|
||||
* definition stores the {@code ref} key and optional {@code args}.
|
||||
*/
|
||||
@FunctionalInterface
|
||||
public interface RoutingCondition {
|
||||
boolean matches(String conditionKey, Map<String, Object> args, ProcessContext context);
|
||||
boolean matches(String conditionKey, Map<String, Object> args, ProcessRuntime runtime);
|
||||
|
||||
static RoutingCondition always() {
|
||||
return (key, args, context) -> true;
|
||||
return (key, args, runtime) -> true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import com.jetlumen.ordo.api.OrdoEventListener;
|
||||
import com.jetlumen.ordo.api.ProcessContext;
|
||||
import com.jetlumen.ordo.api.ProcessDefinition;
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.ProcessRuntime;
|
||||
import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
@@ -368,7 +369,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
|
||||
private boolean escalateReassign(ApprovalTask overdue, ProcessInstance instance, ApprovalStep step, StepDue due,
|
||||
Instant now, List<ProcessEvent> events) {
|
||||
String newAssignee = assigneeResolver.resolve(due.to(), step, instance.context());
|
||||
String newAssignee = assigneeResolver.resolve(due.to(), step, ProcessRuntime.of(instance, step.id()));
|
||||
requireText(newAssignee, "resolved assignee");
|
||||
boolean duplicatePending = taskRepository.findByInstanceIdAndStepId(overdue.instanceId(), overdue.stepId())
|
||||
.stream()
|
||||
@@ -396,8 +397,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
actionExecutionRepository.insert(new ActionExecution(executionId, instance.id(), overdue.stepId(),
|
||||
due.action(), ActionExecutionStatus.PENDING, null,
|
||||
now.plusMillis(eventSequence.getAndIncrement()), null));
|
||||
queued.add(new PendingAction(executionId, due.action(), instance.id(), overdue.stepId(),
|
||||
instance.context()));
|
||||
queued.add(new PendingAction(executionId, due.action(), ProcessRuntime.of(instance, overdue.stepId())));
|
||||
}
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
due.action(), now);
|
||||
@@ -417,7 +417,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
|
||||
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now, List<ProcessEvent> events) {
|
||||
for (String candidate : step.candidates()) {
|
||||
String assignee = assigneeResolver.resolve(candidate, step, instance.context());
|
||||
String assignee = assigneeResolver.resolve(candidate, step, ProcessRuntime.of(instance, step.id()));
|
||||
requireText(assignee, "resolved assignee");
|
||||
Instant dueAt = step.due() == null ? null : now.plus(step.due().after());
|
||||
ApprovalTask task = new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee,
|
||||
@@ -525,8 +525,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
actionExecutionRepository.insert(new ActionExecution(executionId, instance.id(), current.id(),
|
||||
current.actionKey(), ActionExecutionStatus.PENDING, null,
|
||||
now.plusMillis(eventSequence.getAndIncrement()), null));
|
||||
queued.add(new PendingAction(executionId, current.actionKey(), instance.id(), current.id(),
|
||||
instance.context()));
|
||||
queued.add(new PendingAction(executionId, current.actionKey(), ProcessRuntime.of(instance, current.id())));
|
||||
touchToken(definition, instance.id(), current.id());
|
||||
StepTransition matched = resolveTransition(definition, current, instance);
|
||||
if (matched.toStepId() == null) {
|
||||
@@ -588,15 +587,15 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
for (PendingAction pending : queued) {
|
||||
Instant now = clock.instant();
|
||||
try {
|
||||
actionHandler.execute(pending.actionKey(), pending.context());
|
||||
actionHandler.execute(pending.actionKey(), pending.runtime());
|
||||
actionExecutionRepository.complete(pending.executionId(), ActionExecutionStatus.SUCCESS, null, now);
|
||||
record(events, pending.instanceId(), null, pending.stepId(), ProcessEventType.ACTION_SUCCEEDED, null,
|
||||
pending.actionKey(), now);
|
||||
record(events, pending.runtime().instanceId(), null, pending.runtime().stepId(),
|
||||
ProcessEventType.ACTION_SUCCEEDED, null, pending.actionKey(), now);
|
||||
} catch (RuntimeException e) {
|
||||
actionExecutionRepository.complete(pending.executionId(), ActionExecutionStatus.FAILED, e.getMessage(),
|
||||
now);
|
||||
record(events, pending.instanceId(), null, pending.stepId(), ProcessEventType.ACTION_FAILED, null,
|
||||
e.getMessage(), now);
|
||||
record(events, pending.runtime().instanceId(), null, pending.runtime().stepId(),
|
||||
ProcessEventType.ACTION_FAILED, null, e.getMessage(), now);
|
||||
LOG.log(Level.WARNING, "action failed: " + pending.actionKey(), e);
|
||||
}
|
||||
}
|
||||
@@ -632,19 +631,19 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
outgoing.stream().filter(transition -> transition.when() == null).sorted(byPriority).forEach(candidates::add);
|
||||
for (StepTransition transition : candidates) {
|
||||
RoutingWhen when = transition.when();
|
||||
if (when == null || matches(when, instance.context())) {
|
||||
if (when == null || matches(when, instance, step.id())) {
|
||||
return transition;
|
||||
}
|
||||
}
|
||||
throw new NoRouteFoundException(step.id(), instance.id());
|
||||
}
|
||||
|
||||
private boolean matches(RoutingWhen when, ProcessContext context) {
|
||||
private boolean matches(RoutingWhen when, ProcessInstance instance, String stepId) {
|
||||
if (when instanceof RoutingWhen.Predicate predicate) {
|
||||
return RoutingPredicateEvaluator.matches(predicate.tree(), context);
|
||||
return RoutingPredicateEvaluator.matches(predicate.tree(), instance.context());
|
||||
}
|
||||
RoutingWhen.Ref ref = (RoutingWhen.Ref) when;
|
||||
return routingCondition.matches(ref.key(), ref.args(), context);
|
||||
return routingCondition.matches(ref.key(), ref.args(), ProcessRuntime.of(instance, stepId));
|
||||
}
|
||||
|
||||
/** Marks any still-pending sibling candidate tasks for the same step as skipped. */
|
||||
@@ -721,7 +720,6 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
return UUID.randomUUID().toString();
|
||||
}
|
||||
|
||||
private record PendingAction(String executionId, String actionKey, String instanceId, String stepId,
|
||||
ProcessContext context) {
|
||||
private record PendingAction(String executionId, String actionKey, ProcessRuntime runtime) {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import com.jetlumen.ordo.api.ProcessDefinition;
|
||||
import com.jetlumen.ordo.api.ProcessEvent;
|
||||
import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessRuntime;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
import com.jetlumen.ordo.api.RoutingPredicate;
|
||||
@@ -353,8 +354,8 @@ class InMemoryOrdoEngineTest {
|
||||
|
||||
@Test
|
||||
void resolvesAssigneesFromTheProcessContext() {
|
||||
InMemoryOrdoEngine contextAwareEngine = new InMemoryOrdoEngine((AssigneeResolver) (candidate, step, context) ->
|
||||
context.value(step.id())
|
||||
InMemoryOrdoEngine contextAwareEngine = new InMemoryOrdoEngine((AssigneeResolver) (candidate, step, runtime) ->
|
||||
runtime.context().value(step.id())
|
||||
.filter(String.class::isInstance)
|
||||
.map(String.class::cast)
|
||||
.orElse(candidate));
|
||||
@@ -513,12 +514,12 @@ class InMemoryOrdoEngineTest {
|
||||
|
||||
@Test
|
||||
void routesUsingParameterizedHostCondition() {
|
||||
RoutingCondition routingCondition = (key, args, context) -> {
|
||||
RoutingCondition routingCondition = (key, args, runtime) -> {
|
||||
if (!"amountGt".equals(key)) {
|
||||
return false;
|
||||
}
|
||||
Number threshold = (Number) args.get("threshold");
|
||||
return context.value("amount")
|
||||
return runtime.context().value("amount")
|
||||
.filter(Number.class::isInstance)
|
||||
.map(Number.class::cast)
|
||||
.map(amount -> amount.doubleValue() > threshold.doubleValue())
|
||||
@@ -541,8 +542,12 @@ class InMemoryOrdoEngineTest {
|
||||
@Test
|
||||
void runsActionStepsAfterApprovalThenCreatesTheNextApprovalTask() {
|
||||
List<String> executed = new java.util.ArrayList<>();
|
||||
List<ProcessRuntime> runtimes = new java.util.ArrayList<>();
|
||||
InMemoryOrdoEngine actionEngine = new InMemoryOrdoEngine(Clock.systemUTC(), AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> executed.add(key));
|
||||
RoutingCondition.always(), (key, runtime) -> {
|
||||
executed.add(key);
|
||||
runtimes.add(runtime);
|
||||
});
|
||||
actionEngine.publish(new ProcessDefinition("leave", "Leave request", List.of(
|
||||
ApprovalStep.single("manager", "Manager approval", "maria"),
|
||||
ApprovalStep.action("notify", "Notify HR", "leave-approved-mail"),
|
||||
@@ -556,6 +561,12 @@ class InMemoryOrdoEngineTest {
|
||||
actionEngine.approve(actionEngine.findPendingTasksByInstanceId(instance.id()).get(0).id(), "maria");
|
||||
|
||||
assertEquals(List.of("leave-approved-mail"), executed);
|
||||
var runtime = runtimes.get(0);
|
||||
assertEquals(instance.id(), runtime.instanceId());
|
||||
assertEquals("leave", runtime.definitionId());
|
||||
assertEquals(1, runtime.definitionVersion());
|
||||
assertEquals("alice", runtime.initiator());
|
||||
assertEquals("notify", runtime.stepId());
|
||||
ApprovalTask hrTask = actionEngine.findPendingTasksByInstanceId(instance.id()).get(0);
|
||||
assertEquals("hr", hrTask.stepId());
|
||||
assertEquals(ProcessStatus.RUNNING, actionEngine.findInstance(instance.id()).orElseThrow().status());
|
||||
|
||||
Reference in New Issue
Block a user