面向 AI 代理的工作者: 必需子技能:使用 superpowers:subagent-driven-development 逐任务实现此计划。步骤使用复选框(
- [ ])语法来跟踪进度。
目标: 在现有 tasks 表之上新增任务编排模块(含环 DAG:多前置 AND/OR、超时/失败兜底、计数/条件循环终止),后端提供 Flow/Execution 管理与 Quartz 轮询,cfc-web 管理端提供 jsPlumb 可视化编排画布。
架构: 新增 4 张表(flows / edges / executions / node_instances),核心编排引擎 OrchestrationEngine 负责节点触发、循环终止、兜底分支,通过 OrchestrationPollingJob(Quartz)轮询评估,节点执行复用 TaskService 创建 tasks 实例。TaskService 完成/失败点回调编排引擎(仅编排节点任务生效)。
技术栈: Spring Boot 2.7.18 + MyBatis-Plus 3.5.3.1 + Java 8 + Quartz(新增依赖);cfc-web Vue 2 + Element UI + jsPlumb。
设计规格: docs/superpowers/specs/2026-09-19-task-orchestration-design.md
cfc-backend/src/main/java/com/etotem/cfc/entity/TaskOrchestrationFlow.java — flows 表实体cfc-backend/src/main/java/com/etotem/cfc/entity/TaskOrchestrationEdge.java — edges 表实体cfc-backend/src/main/java/com/etotem/cfc/entity/TaskOrchestrationExecution.java — executions 表实体cfc-backend/src/main/java/com/etotem/cfc/entity/TaskOrchestrationNodeInstance.java — node_instances 表实体cfc-backend/src/main/java/com/etotem/cfc/mapper/TaskOrchestrationFlowMapper.javacfc-backend/src/main/java/com/etotem/cfc/mapper/TaskOrchestrationEdgeMapper.javacfc-backend/src/main/java/com/etotem/cfc/mapper/TaskOrchestrationExecutionMapper.javacfc-backend/src/main/java/com/etotem/cfc/mapper/TaskOrchestrationNodeInstanceMapper.javacfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java — 编排引擎(核心)cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationFlowService.java — Flow CRUDcfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationExecutionService.java — Execution 生命周期cfc-backend/src/main/java/com/etotem/cfc/controller/OrchestrationController.java — REST 接口cfc-backend/src/main/java/com/etotem/cfc/task/OrchestrationPollingJob.java — Quartz 轮询 Jobcfc-backend/src/main/java/com/etotem/cfc/config/QuartzConfig.java — Quartz 调度配置cfc-backend/pom.xml — 新增 spring-boot-starter-quartz 依赖cfc-backend/src/main/java/com/etotem/cfc/config/DatabaseInitializer.java — 新增 4 张表迁移cfc-backend/src/main/resources/schema.sql — 同步 4 张表建表语句cfc-backend/src/main/java/com/etotem/cfc/service/TaskService.java — 完成/失败点回调编排引擎cfc-web/src/views/orchestration/OrchestrationFlow.vue — 流列表 + 入口cfc-web/src/views/orchestration/FlowEditor.vue — jsPlumb 编排画布cfc-web/src/views/orchestration/ExecutionDetail.vue — 执行详情只读页cfc-web/src/router/index.js — 注册 3 个路由cfc-web/src/api/orchestration.js(新建)— 接口封装cfc-web/package.json — 新增 jsplumb 依赖cfc-backend/src/test/java/com/etotem/cfc/orchestration/OrchestrationEngineTest.java — 引擎单元测试文件:
TaskOrchestrationFlow.java / TaskOrchestrationEdge.java / TaskOrchestrationExecution.java / TaskOrchestrationNodeInstance.java)DatabaseInitializer.java修改:schema.sql
[ ] 步骤 1:创建 4 个实体类
按 Task.java 现有模式(@Data + @TableName + @TableId(type = IdType.AUTO) + implements Serializable)。字段严格对应设计文档 §3:
TaskOrchestrationFlow.java(@TableName("task_orchestration_flows")):
package com.etotem.cfc.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
@Data
@TableName("task_orchestration_flows")
public class TaskOrchestrationFlow implements Serializable {
@TableId(type = IdType.AUTO)
private Long id;
private String name;
private String description;
private Long creatorId;
private Long familyId;
private Integer version;
private String status;
private String scheduleCron;
private String configJson;
private Date createdAt;
private Date updatedAt;
}
TaskOrchestrationEdge.java(@TableName("task_orchestration_edges")):
package com.etotem.cfc.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
@Data
@TableName("task_orchestration_edges")
public class TaskOrchestrationEdge implements Serializable {
@TableId(type = IdType.AUTO)
private Long id;
private Long flowId;
private String fromNodeId;
private String toNodeId;
private String edgeType;
private String operator;
private Integer sortOrder;
private Date createdAt;
}
TaskOrchestrationExecution.java(@TableName("task_orchestration_executions")):
package com.etotem.cfc.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
@Data
@TableName("task_orchestration_executions")
public class TaskOrchestrationExecution implements Serializable {
@TableId(type = IdType.AUTO)
private Long id;
private Long flowId;
private Integer flowVersion;
private Long familyId;
private Long familyMemberId;
private String triggerSource;
private String status;
private Date startedAt;
private Date finishedAt;
private String errorReason;
private Date createdAt;
private Date updatedAt;
}
TaskOrchestrationNodeInstance.java(@TableName("task_orchestration_node_instances")):
package com.etotem.cfc.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
@Data
@TableName("task_orchestration_node_instances")
public class TaskOrchestrationNodeInstance implements Serializable {
@TableId(type = IdType.AUTO)
private Long id;
private Long executionId;
private String nodeId;
private Long taskId;
private String status;
private Integer generation;
private Date conditionMetAt;
private String loopTerminationReason;
private Date startedAt;
private Date completedAt;
private String errorReason;
private Date createdAt;
private Date updatedAt;
}
每个 Mapper 继承 BaseMapper<T>,模式参考 TaskMapper:
package com.etotem.cfc.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.etotem.cfc.entity.TaskOrchestrationFlow;
public interface TaskOrchestrationFlowMapper extends BaseMapper<TaskOrchestrationFlow> {
}
(其余 3 个 Mapper 同理:TaskOrchestrationEdgeMapper / TaskOrchestrationExecutionMapper / TaskOrchestrationNodeInstanceMapper)
在 runMigrations() 末尾(搜索 // 迁移 找到最新编号,递增),加入 4 张表的 CREATE TABLE IF NOT EXISTS 语句,内容与设计文档 §11 DDL 完全一致(含 condition_met_at 列和唯一索引 uk_node_instance)。
在 schema.sql 末尾追加相同的 4 张表 CREATE TABLE IF NOT EXISTS 语句(与迁移脚本保持一致)。
运行:cd cfc-backend && mvn clean compile
预期:BUILD SUCCESS,无编译错误。
[ ] 步骤 6:Commit
git add cfc-backend/src/main/java/com/etotem/cfc/entity/TaskOrchestration*.java \
cfc-backend/src/main/java/com/etotem/cfc/mapper/TaskOrchestration*Mapper.java \
cfc-backend/src/main/java/com/etotem/cfc/config/DatabaseInitializer.java \
cfc-backend/src/main/resources/schema.sql
git commit -m "feat: 任务编排 4 张新表实体/Mapper/迁移/schema"
文件:
cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java创建:cfc-backend/src/test/java/com/etotem/cfc/orchestration/OrchestrationEngineTest.java
[ ] 步骤 1:编写失败的单元测试(AND/OR 语义 + 循环终止 + 幂等)
OrchestrationEngineTest.java:
package com.etotem.cfc.orchestration;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.etotem.cfc.entity.TaskOrchestrationEdge;
import com.etotem.cfc.entity.TaskOrchestrationNodeInstance;
import com.etotem.cfc.service.OrchestrationEngine;
import org.junit.jupiter.api.Test;
import java.util.Arrays;
import java.util.List;
import static org.junit.jupiter.api.Assertions.*;
public class OrchestrationEngineTest {
private OrchestrationEngine engine = new OrchestrationEngine();
// 前置 AND:两个源节点都 completed 才触发
@Test
public void testAndGateRequiresAllCompleted() {
List<TaskOrchestrationEdge> inEdges = Arrays.asList(
edge("n1", "n3", "AND"),
edge("n2", "n3", "AND")
);
List<TaskOrchestrationNodeInstance> nodes = Arrays.asList(
node("n1", "completed"),
node("n2", "pending")
);
assertFalse(engine.isUnlocked("n3", inEdges, nodes));
}
// 前置 OR:任一源节点 completed 即触发
@Test
public void testOrGateUnlocksWhenAnyCompleted() {
List<TaskOrchestrationEdge> inEdges = Arrays.asList(
edge("n1", "n3", "OR"),
edge("n2", "n3", "OR")
);
List<TaskOrchestrationNodeInstance> nodes = Arrays.asList(
node("n1", "completed"),
node("n2", "pending")
);
assertTrue(engine.isUnlocked("n3", inEdges, nodes));
}
// 循环终止:generation 达到 max_loops 时置 terminated
@Test
public void testLoopTerminationByCount() {
JSONObject cfg = JSON.parseObject("{\"max_loops\":3}");
TaskOrchestrationNodeInstance ni = node("n1", "in_progress");
ni.setGeneration(3); // 已达上限
engine.applyLoopTermination(ni, cfg, 3);
assertEquals("terminated", ni.getStatus());
assertEquals("max_loops_reached", ni.getLoopTerminationReason());
}
private TaskOrchestrationEdge edge(String from, String to, String op) {
TaskOrchestrationEdge e = new TaskOrchestrationEdge();
e.setFromNodeId(from);
e.setToNodeId(to);
e.setOperator(op);
return e;
}
private TaskOrchestrationNodeInstance node(String id, String status) {
TaskOrchestrationNodeInstance n = new TaskOrchestrationNodeInstance();
n.setNodeId(id);
n.setStatus(status);
return n;
}
}
说明:
isUnlocked/applyLoopTermination是引擎将要暴露的纯函数方法,便于单元测试(不依赖数据库)。
运行:cd cfc-backend && mvn test -Dtest=OrchestrationEngineTest
预期:编译失败,OrchestrationEngine 类不存在。
OrchestrationEngine.java 核心方法(注入 4 个 Mapper + TaskService + TaskTemplateMapper,@Resource DI):
package com.etotem.cfc.service;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.etotem.cfc.entity.*;
import com.etotem.cfc.mapper.*;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.*;
import java.util.stream.Collectors;
@Slf4j
@Service
public class OrchestrationEngine {
@Resource
private TaskOrchestrationNodeInstanceMapper nodeInstanceMapper;
@Resource
private TaskOrchestrationEdgeMapper edgeMapper;
@Resource
private TaskOrchestrationExecutionMapper executionMapper;
@Resource
private TaskOrchestrationFlowMapper flowMapper;
@Resource
private TaskTemplateMapper taskTemplateMapper;
@Resource
private TaskService taskService;
/**
* 纯函数:判断目标节点是否满足入边触发条件。
* @return true = 可触发
*/
public boolean isUnlocked(String targetNodeId,
List<TaskOrchestrationEdge> inEdges,
List<TaskOrchestrationNodeInstance> nodes) {
if (inEdges == null || inEdges.isEmpty()) {
return true; // 无前置,start_node 直接触发
}
Map<String, String> statusById = nodes.stream()
.collect(Collectors.toMap(TaskOrchestrationNodeInstance::getNodeId,
TaskOrchestrationNodeInstance::getStatus));
// 按 operator 分组:AND 组内全 completed,OR 组内任一 completed
boolean hasAnd = false;
boolean hasOr = false;
boolean orSatisfied = false;
for (TaskOrchestrationEdge e : inEdges) {
String src = e.getFromNodeId();
String srcStatus = statusById.getOrDefault(src, "pending");
if ("AND".equalsIgnoreCase(e.getOperator())) {
hasAnd = true;
if (!"completed".equals(srcStatus)) {
return false; // 任一 AND 源未完成即不触发
}
} else { // OR
hasOr = true;
if ("completed".equals(srcStatus)) {
orSatisfied = true;
}
}
}
if (hasAnd && !hasOr) {
return true; // 纯 AND,全部完成(前面已通过 return false 检查)
}
if (!hasAnd && hasOr) {
return orSatisfied; // 纯 OR
}
// AND + OR 混合:AND 全完成(已通过检查)且 OR 至少一个完成
return true; // 到达此处说明 AND 已全完成,OR 是否有满足不影响(按 AND 优先语义)
}
/**
* 纯函数:应用循环终止规则。
*/
public void applyLoopTermination(TaskOrchestrationNodeInstance ni,
JSONObject nodeCfg,
int maxLoops) {
if (ni.getGeneration() != null && ni.getGeneration() >= maxLoops) {
ni.setStatus("terminated");
ni.setLoopTerminationReason("max_loops_reached");
}
}
/**
* 事件驱动评估:节点完成/失败/超时后调用,扫描 pending 节点并触发满足条件的。
*/
public void evaluateFlow(Long executionId) {
// 1. 加载执行实例
TaskOrchestrationExecution exec = executionMapper.selectById(executionId);
if (exec == null || !"running".equals(exec.getStatus())) {
return;
}
// 2. 加载该 execution 的所有节点实例
List<TaskOrchestrationNodeInstance> nodes = nodeInstanceMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, executionId));
// 3. 加载 flow 的边(此处需从 execution.flowId 关联 edges,但版本快照简化:直接读当前边表)
List<TaskOrchestrationEdge> edges = edgeMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationEdge>()
.eq(TaskOrchestrationEdge::getFlowId, exec.getFlowId()));
// 4. 对每个 pending 节点,检查入边是否满足
Set<String> nodeIds = nodes.stream()
.map(TaskOrchestrationNodeInstance::getNodeId).collect(Collectors.toSet());
for (String nodeId : nodeIds) {
List<TaskOrchestrationNodeInstance> pendingNodes = nodes.stream()
.filter(n -> "pending".equals(n.getStatus()) && nodeId.equals(n.getNodeId()))
.collect(Collectors.toList());
if (pendingNodes.isEmpty()) continue;
List<TaskOrchestrationEdge> inEdges = edges.stream()
.filter(e -> nodeId.equals(e.getToNodeId()))
.collect(Collectors.toList());
if (isUnlocked(nodeId, inEdges, nodes)) {
triggerNode(exec, pendingNodes.get(0), inEdges);
}
}
// 5. 判断 flow 是否完成(所有节点均非 pending/running/in_progress)
boolean allTerminal = nodes.stream().allMatch(n ->
!"pending".equals(n.getStatus())
&& !"in_progress".equals(n.getStatus())
&& !"started".equals(n.getStatus()));
if (allTerminal) {
exec.setStatus("completed");
exec.setFinishedAt(new Date());
executionMapper.updateById(exec);
}
}
private void triggerNode(TaskOrchestrationExecution exec,
TaskOrchestrationNodeInstance ni,
List<TaskOrchestrationEdge> inEdges) {
// 幂等:唯一索引兜底,这里先检查
Long existingCount = nodeInstanceMapper.selectCount(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, ni.getExecutionId())
.eq(TaskOrchestrationNodeInstance::getNodeId, ni.getNodeId())
.eq(TaskOrchestrationNodeInstance::getGeneration, ni.getGeneration())
.ne(TaskOrchestrationNodeInstance::getStatus, "pending"));
if (existingCount != null && existingCount > 0) return;
// triggerNode 完整实现见任务 3(步骤 1),此处暂不展开,
// 任务 2 仅保留 evaluateFlow 骨架及 isUnlocked / applyLoopTermination 纯函数
}
// triggerNode、createTaskFromTemplate、triggerFailedEdges 等完整实现见任务 3(步骤 1)
}
注意:任务 2 仅定义引擎骨架(
evaluateFlow、isUnlocked、applyLoopTermination、onTaskCompleted、onTaskFailed、onTaskTimeout、findByTaskId)。triggerNode的完整实现(含createTaskFromTemplate模板字段映射)见任务 3 步骤 1。
运行:cd cfc-backend && mvn test -Dtest=OrchestrationEngineTest
预期:PASS(4 个测试全部通过)。
运行:cd cfc-backend && mvn clean compile
预期:BUILD SUCCESS。
[ ] 步骤 6:Commit
git add cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java \
cfc-backend/src/test/java/com/etotem/cfc/orchestration/OrchestrationEngineTest.java
git commit -m "feat: 编排引擎核心(AND/OR 触发 + 循环终止 + evaluateFlow)"
文件:
cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java — 补全 triggerNode修改:cfc-backend/src/main/java/com/etotem/cfc/service/TaskService.java — 完成/失败点回调
[ ] 步骤 1:补全 triggerNode(加载模板 → 创建 tasks → 回写)
在 OrchestrationEngine 中补全 triggerNode 方法体,替换步骤 2 中的 TODO 部分。核心逻辑:
private void triggerNode(TaskOrchestrationExecution exec,
TaskOrchestrationNodeInstance ni,
List<TaskOrchestrationEdge> inEdges) {
// 幂等检查(唯一索引兜底)
Long existingCount = nodeInstanceMapper.selectCount(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, ni.getExecutionId())
.eq(TaskOrchestrationNodeInstance::getNodeId, ni.getNodeId())
.eq(TaskOrchestrationNodeInstance::getGeneration, ni.getGeneration())
.ne(TaskOrchestrationNodeInstance::getStatus, "pending"));
if (existingCount != null && existingCount > 0) return;
// 1. 读 flow 的 config_json,拿到节点配置
TaskOrchestrationFlow flow = flowMapper.selectById(exec.getFlowId());
JSONObject flowCfg = JSON.parseObject(flow.getConfigJson());
JSONObject nodeCfg = null;
for (Object o : flowCfg.getJSONArray("nodes")) {
JSONObject n = (JSONObject) o;
if (ni.getNodeId().equals(n.getString("id"))) {
nodeCfg = n;
break;
}
}
if (nodeCfg == null) {
ni.setStatus("failed");
ni.setErrorReason("节点配置不存在");
nodeInstanceMapper.updateById(ni);
return;
}
// 2. 加载任务模板
String templateRef = nodeCfg.getString("task_template_ref");
TaskTemplate template = taskTemplateMapper.selectById(Long.valueOf(templateRef));
if (template == null) {
ni.setStatus("failed");
ni.setErrorReason("任务模板不存在: " + templateRef);
nodeInstanceMapper.updateById(ni);
triggerFailedEdges(exec, ni); // 触发 failed 边
return;
}
// 3. 创建 tasks 实例(调用 TaskService.createTask(Long, CreateTaskDTO))
// 根据模板配置填充 DTO,关键字段映射:
// title/description:从 template 读取
// points:从 template.pts 读取
// deadline:nodeCfg.getInteger("timeout_minutes") 对应 task.deadline
// is_daily_progress:nodeCfg.boolean("is_daily_progress")
// source_type/source_id:flow config_json 中的对应字段
// memberOnly:template.integer("memberOnly")
CreateTaskDTO dto = new CreateTaskDTO();
dto.setTitle(template.getTitle());
dto.setDescription(template.getDescription());
dto.setPoints(template.getPoints());
dto.setTaskType(template.getTaskType());
dto.setFrequency(template.getFrequency());
dto.setMaxFrequency(template.getMaxFrequency());
dto.setNeedReview(template.getNeedReview());
dto.setReviewType(template.getReviewType());
dto.setReviewByCategory(template.getReviewByCategory());
dto.setCompleteTypes(template.getCompleteTypes());
dto.setEarliestStart(template.getEarliestStart());
dto.setLatestEnd(template.getLatestEnd());
dto.setExecutorType(nodeCfg.getString("executorType") != null ? nodeCfg.getString("executorType") : "child");
dto.setExecutorId(ni.getFamilyMemberId());
dto.setMemberId(ni.getFamilyMemberId());
dto.setMinigameCode(nodeCfg.getString("minigameCode"));
dto.setRepeatType(nodeCfg.getString("repeatType"));
dto.setDuration(nodeCfg.getInteger("duration"));
dto.setDimensionCode(nodeCfg.getString("dimensionCode"));
dto.setDimensionWeights(nodeCfg.getString("dimensionWeights"));
dto.setMemberOnly(nodeCfg.getInteger("memberOnly") != null ? nodeCfg.getInteger("memberOnly") : 0);
dto.setPrerequisiteTaskId(nodeCfg.getLong("prerequisiteTaskId"));
dto.setActionType(nodeCfg.getString("actionType"));
dto.setActionConfig(nodeCfg.getString("actionConfig"));
dto.setRequireInput(nodeCfg.getInteger("requireInput"));
dto.setStartRequired(nodeCfg.getInteger("startRequired"));
dto.setMinDurationSeconds(nodeCfg.getInteger("minDurationSeconds"));
dto.setSourceType(nodeCfg.getString("sourceType"));
dto.setSourceId(nodeCfg.getLong("sourceId"));
dto.setIsDailyProgress(nodeCfg.getInteger("is_daily_progress"));
dto.setTargetValue(nodeCfg.getInteger("targetValue"));
// 超时时间:deadline = now + timeout_minutes(若配置)
Long taskId = taskService.createTask(ni.getFamilyMemberId(), dto);
// 4. 回写 task_id + 状态
ni.setTaskId(taskId);
ni.setStatus("in_progress");
ni.setStartedAt(new Date());
nodeInstanceMapper.updateById(ni);
}
private void triggerFailedEdges(TaskOrchestrationExecution exec,
TaskOrchestrationNodeInstance ni) {
// 触发所有 edge_type=failed 且 from_node_id=ni.nodeId 的出边目标节点(创建 pending 实例)
List<TaskOrchestrationEdge> failedEdges = edgeMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationEdge>()
.eq(TaskOrchestrationEdge::getFlowId, exec.getFlowId())
.eq(TaskOrchestrationEdge::getFromNodeId, ni.getNodeId())
.eq(TaskOrchestrationEdge::getEdgeType, "failed"));
for (TaskOrchestrationEdge e : failedEdges) {
// 幂等:检查同execution_id + to_node_id + generation=0 的 pending 实例是否存在
Long exists = nodeInstanceMapper.selectCount(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, exec.getId())
.eq(TaskOrchestrationNodeInstance::getNodeId, e.getToNodeId())
.eq(TaskOrchestrationNodeInstance::getGeneration, 0)
.eq(TaskOrchestrationNodeInstance::getStatus, "pending"));
if (exists != null && exists > 0) continue; // 已存在则跳过
TaskOrchestrationNodeInstance newNi = new TaskOrchestrationNodeInstance();
newNi.setExecutionId(exec.getId());
newNi.setNodeId(e.getToNodeId());
newNi.setGeneration(0);
newNi.setStatus("pending");
newNi.setCreatedAt(new Date());
nodeInstanceMapper.insert(newNi);
}
}
关键:
triggerFailedEdges通过查询唯一索引(execution_id, node_id, generation)实现幂等,避免 Quartz 轮询与回调同时触发导致重复创建节点任务。createTaskFromTemplate内部需调用TaskService.createTask(CreateTaskDTO),字段映射(title/description/points/needReview/deadline)参考TaskService现有模板实例化逻辑。实施时先读取TaskService.createTask方法签名确认参数,不可臆测。
定位 TaskService 中任务完成(状态置 completed)的代码点,末尾添加(仅编排节点任务生效):
@Resource
@Lazy
private OrchestrationEngine orchestrationEngine;
// 在任务完成逻辑后:
orchestrationEngine.onTaskCompleted(taskId);
在 OrchestrationEngine 新增回调方法:
public void onTaskCompleted(Long taskId) {
TaskOrchestrationNodeInstance ni = findByTaskId(taskId);
if (ni == null) return; // 非编排节点任务,零开销
ni.setStatus("completed");
ni.setCompletedAt(new Date());
nodeInstanceMapper.updateById(ni);
evaluateFlow(ni.getExecutionId());
}
private TaskOrchestrationNodeInstance findByTaskId(Long taskId) {
return nodeInstanceMapper.selectOne(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getTaskId, taskId));
}
同理,在任务失败/放弃逻辑点调用 orchestrationEngine.onTaskFailed(taskId),引擎实现:
public void onTaskFailed(Long taskId) {
TaskOrchestrationNodeInstance ni = findByTaskId(taskId);
if (ni == null) return;
ni.setStatus("failed");
ni.setCompletedAt(new Date());
nodeInstanceMapper.updateById(ni);
// 触发 failed 边
TaskOrchestrationExecution exec = executionMapper.selectById(ni.getExecutionId());
triggerFailedEdges(exec, ni);
evaluateFlow(ni.getExecutionId());
}
运行:cd cfc-backend && mvn clean compile
预期:BUILD SUCCESS(确认 TaskService.createTask 签名匹配,flowMapper 已注入)。
[ ] 步骤 5:Commit
git add cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java \
cfc-backend/src/main/java/com/etotem/cfc/service/TaskService.java
git commit -m "feat: 节点任务创建 + TaskService 完成/失败回调编排引擎"
文件:
cfc-backend/pom.xml — 新增 Quartz 依赖cfc-backend/src/main/java/com/etotem/cfc/config/QuartzConfig.javacfc-backend/src/main/java/com/etotem/cfc/task/OrchestrationPollingJob.java修改:cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java — 新增 onTaskTimeout
[ ] 步骤 1:pom.xml 新增 Quartz 依赖
在 <dependencies> 中加入:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-quartz</artifactId>
</dependency>
[ ] 步骤 2:创建 QuartzConfig
package com.etotem.cfc.config;
import org.quartz.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class QuartzConfig {
@Bean
public JobDetail orchestrationPollingJobDetail() {
return JobBuilder.newJob(OrchestrationPollingJob.class)
.withIdentity("orchestrationPollingJob")
.storeDurably()
.build();
}
@Bean
public Trigger orchestrationPollingTrigger() {
return TriggerBuilder.newTrigger()
.forJob(orchestrationPollingJobDetail())
.withIdentity("orchestrationPollingTrigger")
.withSchedule(SimpleScheduleBuilder.simpleSchedule()
.withIntervalInSeconds(30)
.repeatForever())
.build();
}
}
注意:
OrchestrationPollingJob在com.etotem.cfc.task包,QuartzConfig在com.etotem.cfc.config包,需 import 对应类。
[ ] 步骤 3:创建 OrchestrationPollingJob
package com.etotem.cfc.task;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.etotem.cfc.entity.Task;
import com.etotem.cfc.entity.TaskOrchestrationExecution;
import com.etotem.cfc.entity.TaskOrchestrationNodeInstance;
import com.etotem.cfc.mapper.TaskMapper;
import com.etotem.cfc.mapper.TaskOrchestrationExecutionMapper;
import com.etotem.cfc.mapper.TaskOrchestrationNodeInstanceMapper;
import com.etotem.cfc.service.OrchestrationEngine;
import lombok.extern.slf4j.Slf4j;
import org.quartz.DisallowConcurrentExecution;
import org.quartz.Job;
import org.quartz.JobExecutionContext;
import javax.annotation.Resource;
import java.util.Date;
import java.util.List;
@Slf4j
@DisallowConcurrentExecution
public class OrchestrationPollingJob implements Job {
@Resource
private TaskOrchestrationExecutionMapper executionMapper;
@Resource
private TaskOrchestrationNodeInstanceMapper nodeInstanceMapper;
@Resource
private TaskMapper taskMapper;
@Resource
private OrchestrationEngine orchestrationEngine;
@Override
public void execute(JobExecutionContext context) {
log.info("编排轮询开始...");
try {
// 1. 查询所有 running 状态的 execution(每次最多 20 个)
List<TaskOrchestrationExecution> executions = executionMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationExecution>()
.eq(TaskOrchestrationExecution::getStatus, "running")
.last("LIMIT 20"));
for (TaskOrchestrationExecution exec : executions) {
// 2. 检查 timeout 节点(deadline 到期)
checkTimeouts(exec);
// 3. 评估 pending 节点
orchestrationEngine.evaluateFlow(exec.getId());
}
} catch (Exception e) {
log.error("编排轮询异常", e);
}
log.info("编排轮询结束");
}
private void checkTimeouts(TaskOrchestrationExecution exec) {
List<TaskOrchestrationNodeInstance> nodes = nodeInstanceMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, exec.getId())
.eq(TaskOrchestrationNodeInstance::getStatus, "in_progress"));
Date now = new Date();
for (TaskOrchestrationNodeInstance ni : nodes) {
if (ni.getTaskId() == null) continue;
Task task = taskMapper.selectById(ni.getTaskId());
if (task != null && task.getDeadline() != null && task.getDeadline().before(now)
&& !"completed".equals(task.getStatus())) {
orchestrationEngine.onTaskTimeout(ni.getTaskId());
}
}
}
}
注意:Quartz 的
Job实例由 Quartz 容器管理,@Resource注入需要SpringBeanJobFactory支持。若注入失败,需在QuartzConfig配置SpringBeanJobFactory(让 Quartz 使用 Spring 容器创建 Job 实例)。实施时优先验证@Resource是否生效,不行则加SpringBeanJobFactorybean。
[ ] 步骤 4:OrchestrationEngine 新增 onTaskTimeout
public void onTaskTimeout(Long taskId) {
TaskOrchestrationNodeInstance ni = findByTaskId(taskId);
if (ni == null) return;
ni.setStatus("timeout");
ni.setCompletedAt(new Date());
nodeInstanceMapper.updateById(ni);
TaskOrchestrationExecution exec = executionMapper.selectById(ni.getExecutionId());
// 触发 timeout 边(兜底)
List<TaskOrchestrationEdge> timeoutEdges = edgeMapper.selectList(
new LambdaQueryWrapper<TaskOrchestrationEdge>()
.eq(TaskOrchestrationEdge::getFlowId, exec.getFlowId())
.eq(TaskOrchestrationEdge::getFromNodeId, ni.getNodeId())
.eq(TaskOrchestrationEdge::getEdgeType, "timeout"));
for (TaskOrchestrationEdge e : timeoutEdges) {
// 幂等:检查同execution_id + to_node_id + generation=0 的 pending 实例是否存在
Long exists = nodeInstanceMapper.selectCount(
new LambdaQueryWrapper<TaskOrchestrationNodeInstance>()
.eq(TaskOrchestrationNodeInstance::getExecutionId, exec.getId())
.eq(TaskOrchestrationNodeInstance::getNodeId, e.getToNodeId())
.eq(TaskOrchestrationNodeInstance::getGeneration, 0)
.eq(TaskOrchestrationNodeInstance::getStatus, "pending"));
if (exists != null && exists > 0) continue; // 已存在则跳过
TaskOrchestrationNodeInstance newNi = new TaskOrchestrationNodeInstance();
newNi.setExecutionId(exec.getId());
newNi.setNodeId(e.getToNodeId());
newNi.setGeneration(0);
newNi.setStatus("pending");
newNi.setCreatedAt(new Date());
nodeInstanceMapper.insert(newNi);
}
evaluateFlow(exec.getId());
}
[ ] 步骤 5:编译验证
运行:cd cfc-backend && mvn clean compile
预期:BUILD SUCCESS。
[ ] 步骤 6:Commit
git add cfc-backend/pom.xml \
cfc-backend/src/main/java/com/etotem/cfc/config/QuartzConfig.java \
cfc-backend/src/main/java/com/etotem/cfc/task/OrchestrationPollingJob.java \
cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationEngine.java
git commit -m "feat: Quartz 轮询 + timeout 检查 + 兜底边触发"
文件:
cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationFlowService.javacfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationExecutionService.java创建:cfc-backend/src/main/java/com/etotem/cfc/controller/OrchestrationController.java
[ ] 步骤 1:创建 OrchestrationFlowService(Flow CRUD)
按设计文档 §7.1 实现 6 个方法:save / publish / list / detail / archive / delete。核心逻辑(save 保存 config_json + edges 快照,publish 状态转换 + version+1)。
按 §7.2 实现:start(创建 execution + start_node 触发)/ pause / resume / terminate / detail / list。start 方法:
@Transactional
public Result<Map<String, Object>> start(Long flowId, Long familyMemberId, Long userId) {
TaskOrchestrationFlow flow = flowMapper.selectById(flowId);
if (flow == null) return Result.error("编排流不存在");
if (!"published".equals(flow.getStatus())) return Result.error("编排流未发布");
// 创建执行实例
TaskOrchestrationExecution exec = new TaskOrchestrationExecution();
exec.setFlowId(flowId);
exec.setFlowVersion(flow.getVersion());
exec.setFamilyMemberId(familyMemberId);
exec.setTriggerSource("manual");
exec.setStatus("running");
exec.setStartedAt(new Date());
executionMapper.insert(exec);
// 读取 config_json,创建 start_node 的 pending 实例
JSONObject cfg = JSON.parseObject(flow.getConfigJson());
for (Object o : cfg.getJSONArray("nodes")) {
JSONObject n = (JSONObject) o;
TaskOrchestrationNodeInstance ni = new TaskOrchestrationNodeInstance();
ni.setExecutionId(exec.getId());
ni.setNodeId(n.getString("id"));
ni.setGeneration(0);
if (n.getBooleanValue("is_start_node")) {
ni.setStatus("pending"); // 交由 evaluateFlow 触发
} else {
ni.setStatus("pending");
}
ni.setCreatedAt(new Date());
nodeInstanceMapper.insert(ni);
}
// 触发 start_node
orchestrationEngine.evaluateFlow(exec.getId());
Map<String, Object> data = new HashMap<>();
data.put("executionId", exec.getId());
return Result.success(data);
}
注意:
familyMemberId对应FamilyMember.id(主键),在start方法中已直接注入,无需额外反查。若需关联家庭属性(如familyId、isAdmin等),可通过FamilyMemberMapper.selectById(familyMemberId)实现。
统一 @PostMapping,路由前缀 /api/orchestration,注入两个 Service + 手动检查角色(@RequestAttribute("role"))。实现 §7.1/§7.2/§7.3 全部接口,含 /node/fail、/node/restart、/node/condition-met。
运行:cd cfc-backend && mvn clean compile
预期:BUILD SUCCESS。
[ ] 步骤 5:Commit
git add cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationFlowService.java \
cfc-backend/src/main/java/com/etotem/cfc/service/OrchestrationExecutionService.java \
cfc-backend/src/main/java/com/etotem/cfc/controller/OrchestrationController.java
git commit -m "feat: Flow/Execution Service + OrchestrationController REST 接口"
文件:
创建:cfc-backend/src/test/java/com/etotem/cfc/orchestration/OrchestrationIntegrationTest.java
[ ] 步骤 1:编写集成测试(完整流执行)
覆盖:start → node1 完成 → node2 触发 → node3 超时 → 兜底 node4。用真实 Mapper(测试库)验证幂等(回调 + 轮询不重复建任务)。
运行:cd cfc-backend && mvn test
预期:PASS(含新增集成测试)。
运行:cd cfc-backend && mvn spring-boot:run(确认无 Unknown column 报错,4 张表迁移成功)。
[ ] 步骤 4:Commit
git add cfc-backend/src/test/java/com/etotem/cfc/orchestration/OrchestrationIntegrationTest.java
git commit -m "test: 编排引擎集成测试(完整流 + 幂等)"
文件:
cfc-web/src/views/orchestration/OrchestrationFlow.vuecfc-web/src/views/orchestration/FlowEditor.vuecfc-web/src/views/orchestration/ExecutionDetail.vuecfc-web/src/api/orchestration.jscfc-web/src/router/index.js修改:cfc-web/package.json — 新增 jsplumb
[ ] 步骤 1:安装 jsplumb + 封装 API
cd cfc-web && npm install jsplumb --save
cfc-web/src/api/orchestration.js:封装 12 个接口(flow 6 个 + execution 6 个)。
Element UI 表格:流名称、版本、状态、创建时间、操作(编辑/发布/归档/删除/启动执行)。参考现有管理端列表页模式。
[ ] 步骤 3:创建 FlowEditor.vue(jsPlumb 画布)
左侧:可拖拽任务模板列表(admin_task_templates 接口)
中央:jsPlumb 画布(节点拖拽 + 连线 + 删除连线)
右侧:属性面板(超时时间、max_loops、终止条件、AND/OR 边类型)
工具栏:保存草稿 / 发布(发布前校验孤立节点、环必须有终止条件、task_template_ref 存在、有且仅一个 start_node)
[ ] 步骤 4:创建 ExecutionDetail.vue(只读树)
树形展示节点实例状态 + 关联 tasks 记录 + 手动干预(标记失败/终止流/重启节点/条件达成)。
cfc-web/src/router/index.js 新增:
{ path: '/orchestration', component: OrchestrationFlow },
{ path: '/orchestration/edit/:id', component: FlowEditor },
{ path: '/orchestration/execution/:id', component: ExecutionDetail }
运行:cd cfc-web && npm run build
预期:构建成功。
[ ] 步骤 7:Commit
git add cfc-web/src/views/orchestration/ cfc-web/src/api/orchestration.js \
cfc-web/src/router/index.js cfc-web/package.json cfc-web/package-lock.json
git commit -m "feat: 任务编排管理端(列表 + jsPlumb 画布 + 执行详情)"
| 不做 | 原因 |
|---|---|
| 小程序端编排编辑 | 编排编辑只放 cfc-web 管理端 |
| 流市场/模板市场 | 后续迭代 |
| 节点级并行执行 | 一个节点同一时刻只一个任务实例 |
| 家庭挑战/五维打卡合并 | 已有设计文档明确不合并 |
| flow 配置的版本快照隔离(edges 版本化) | 当前 edges 表读最新,版本隔离仅对 execution 记录 flow_version;后续如需精确快照再补 edges 版本列 |
mvn clean compile 和 npm run build