文档目的:在现有 DAG 工作流引擎基础上,新增一种"自主智能体"执行模式。在该模式下,画布作为智能体的"能力地图"(参考路径,非约束),由全局 Planner Agent 决定每一步访问哪个节点;同一节点可被多次访问,每次访问独立暂存工作区快照。
改造范围:
- 后端:新增
AgenticExecutor+PlannerService+WorkflowRunVisit实体(独立路径,不动 DAG 引擎)- 前端:编辑器加 mode 切换、运行结果改为按 Visit 链展示
状态:待实施
我做的不是传统的工作流智能体,而是灵活的可自主决定如何运行的智能体,工作流只是给它的参考而非桎梏。所以,可以创造画布上不存在的边,只要决策 AI 觉得有必要。对应的,每个节点,由于可能经过多次,所以要暂存每次运行时的工作区;若下一次运行(从后继节点飞回来),则将后继节点的工作区变量迁移过来,再次记录。相应地,在"运行结果"处,也会多次出现该节点,每次出现时,工作区变量都可能不同。 如果出现灵活运行的边,要在运行结果处给飞到的节点打一个标识,说明它是从哪个节点飞过来的。
| 维度 | DAG 模式(现有) | AGENTIC 模式(本次新增) |
|---|---|---|
| 画布连线 | 必须严格遵守的控制流 | 参考路径,非约束 |
| 节点语义 | 固定步骤,执行一次 | 能力/工具,可被任意次访问 |
| 引擎角色 | 调度器("该你了") | 执行体(决策者另有其人) |
| 决策权 | 静态画布连线 | 全局 Planner Agent(运行时 LLM) |
| 节点输出 | 单次写入命名空间 | 每次 Visit 独立快照 |
| 终止条件 | 所有节点完成 | 访问到 Output 节点 |
| 决策 | 选择 | 原文 |
|---|---|---|
| 决策者 | X2 全局 Planner | "工作流只是给它的参考" 暗示有"它"主体 |
| 上下文迁移 | 链式继承 + 最近覆盖 | 用户给出明确语义:A#2 工作区 = C#1 工作区 + A#2 输出(覆盖 A#1 同名变量) |
| 跳转约束 | 画布节点集合内任意 | W2,平衡灵活性和可调试性 |
| 终止条件 | 必须到 Output | T1,保证有明确产出 |
访问序列:A → B → C → A(飞回)
A#1 工作区 = 初始输入 + A#1 输出
B#1 工作区 = A#1 工作区 + B#1 输出
C#1 工作区 = B#1 工作区 + C#1 输出
A#2 工作区 = C#1 工作区 + A#2 输出(A#1 同名变量被覆盖)
关键不变式:每个 Visit 的工作区 = 上一个 Visit 工作区 + 本节点输出,同名变量最近覆盖。
WorkflowLevelExecutor 严格执行 DAG,禁止环(DagResolver.java:108-110 检测到环抛异常)ConditionExecutor.java 是 LLM 驱动的分支选择(仅在画布已画的出边里选),可作为 Planner prompt 的实现参考NodeOutputEnvelope,节点级命名空间隔离在 WorkflowContext.nodeScopedOutputsSseEventBus(双轨:emitter + 缓存)支持断线重连本次为全新独立路径,不动 DAG 引擎任何代码。
| 编号 | 行为 |
|---|---|
| B1 | 工作流可配置 mode = DAG(默认)/ AGENTIC |
| B2 | AGENTIC 模式下,运行时由全局 Planner Agent 决定下一步访问哪个节点 |
| B3 | Planner 可指定画布上任意节点(不要求画布有连线),也可指定 Output 节点终结 |
| B4 | 每个节点访问生成独立的 Visit 快照,含工作区变量、输出、Planner 决策原因 |
| B5 | 运行结果按 Visit 顺序展开,同节点多次访问逐条显示,飞回的 Visit 标识来源 |
| B6 | 兜底:maxVisits 硬上限 + Planner 失败走画布参考路径 |
| 编号 | 现有行为 | 保留方式 |
|---|---|---|
| S1 | mode = DAG 的所有现有行为 |
完全不动 WorkflowLevelExecutor,分流入口在 WorkflowEngine.execute() |
| S2 | 节点执行器(NodeExecutor 实现)零改动 |
复用同一套 executor,只换调度层 |
| S3 | NodeOutputEnvelope 输出格式 |
每个 Visit 的 output 仍是 envelope |
| S4 | WorkflowRun 持久化 + SSE 推送 |
Visit 表关联 runId,新事件类型并行推送 |
| S5 | 节点输入解析(NodeInputResolver) |
AGENTIC 模式下 workspace.variables 已被填满,解析器照常工作 |
| S6 | SSE 断线重连 | 新事件类型走同一个 SseEventBus |
新建 AgenticExecutor 类,与 WorkflowLevelExecutor 并列,在 WorkflowEngine.execute() 入口按 mode 分流。
优点:
缺点:
破坏单一职责,调试时难以分辨"这段代码是 DAG 还是 AGENTIC 在跑"。
每节点自己决定下一步,违反"全局决策者"语义。
文件:backend/src/main/java/com/agent/management/model/entity/Workflow.java
/** 运行模式:DAG(严格按画布连线执行)/ AGENTIC(智能体自主决策) */
@Column(name = "mode", nullable = false)
private String mode = "DAG";
加 PrePersist 兜底:if (mode == null) mode = "DAG";
新文件:backend/src/main/java/com/agent/management/model/entity/WorkflowRunVisit.java
@Entity
@Table(name = "workflow_run_visits")
@Data
public class WorkflowRunVisit {
@Id @GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "run_id", nullable = false)
private Long runId; // 关联 workflow_runs.id
@Column(name = "visit_uuid", nullable = false, unique = true)
private String visitUuid; // UUID,便于前端引用
@Column(name = "node_id", nullable = false)
private String nodeId; // 访问的节点 id
@Column(name = "node_type")
private String nodeType; // 节点类型快照
@Column(name = "node_label")
private String nodeLabel; // 节点标签快照
@Column(name = "iter_seq", nullable = false)
private Integer iterSeq; // 同节点第几次访问(1-based)
@Column(name = "from_visit_uuid")
private String fromVisitUuid; // 上一个 visit 的 UUID(起始为 null)
@Column(name = "from_node_id")
private String fromNodeId; // 飞回标识:上一个 visit 的 nodeId(起始为 null)
@Lob
@Column(name = "variables_json", columnDefinition = "TEXT")
private String variablesJson; // 本 visit 工作区快照
@Lob
@Column(name = "output_json", columnDefinition = "TEXT")
private String outputJson; // 本节点本次输出 envelope
@Lob
@Column(name = "decision_json", columnDefinition = "TEXT")
private String decisionJson; // Planner 决策原因(next, reason)
@Column(name = "status", nullable = false)
private String status; // RUNNING / SUCCESS / FAILED
@Column(name = "error_msg")
private String errorMsg;
@Column(name = "sort_order", nullable = false)
private Integer sortOrder; // 全局顺序
@Column(name = "created_at")
private LocalDateTime createdAt;
}
新文件:backend/src/main/java/com/agent/management/repository/WorkflowRunVisitRepository.java
public interface WorkflowRunVisitRepository extends JpaRepository<WorkflowRunVisit, Long> {
List<WorkflowRunVisit> findByRunIdOrderBySortOrderAsc(Long runId);
Integer countByRunIdAndNodeId(Long runId, String nodeId);
}
文件:backend/src/main/java/com/agent/management/model/entity/WorkflowRun.java
@Column(name = "visit_count")
private Integer visitCount; // AGENTIC 模式下的总访问次数;DAG 模式为 null
新文件:backend/src/main/java/com/agent/management/engine/AgenticExecutor.java
核心循环(伪代码):
@Slf4j
@Component
@RequiredArgsConstructor
public class AgenticExecutor {
private final Map<String, NodeExecutor> executorMap; // Spring 注入
private final PlannerService plannerService;
private final WorkflowRunVisitRepository visitRepo;
private final NodeWorkspaceBuilder workspaceBuilder;
private final ChatClient.Builder fallbackBuilder;
private final AiModelService aiModelService;
/** 单 run 最多访问 50 次(含 Output 节点),任何情况都强制终止 */
private static final int MAX_VISITS = 50;
/** 同节点连续访问 5 次仍找不到 Output,强制终止 */
private static final int MAX_SAME_NODE_REVISITS = 5;
public void execute(ResolvedDag dag, WorkflowContext ctx, WorkflowRun run,
SseEmitter emitter, SseEventBus eventBus) {
String runId = run.getRunUuid();
Map<String, Integer> nodeVisitCounter = new HashMap<>(); // nodeId -> 已访问次数
Map<String, Object> currentVars = new LinkedHashMap<>(ctx.getInitialInputs());
String currentVisitUuid = null;
String fromNodeId = null;
String fromVisitUuid = null;
int sortOrder = 0;
// 起始节点:userInput 类型优先;否则取入度为 0 的第一个
String currentNodeId = findStartNode(dag);
while (sortOrder < MAX_VISITS) {
// 0. 同节点连续访问上限兜底
int revisits = nodeVisitCounter.merge(currentNodeId, 1, Integer::sum);
if (revisitas > MAX_SAME_NODE_REVISITS) {
abort(ctx, emitter, eventBus, runId,
"节点 " + currentNodeId + " 连续访问超过 " + MAX_SAME_NODE_REVISITS + " 次,疑似死循环");
return;
}
// 1. 推送 nodeRunning SSE
safeSend(emitter, WorkflowRunEvent.nodeRunning(runId, currentNodeId));
// 2. 构建 workspace(把 currentVars 作为变量来源)
NodeWorkspace ws = workspaceBuilder.buildForAgentic(
currentNodeId, dag, ctx, currentVars);
// 3. 调用 executor
NodeExecutionResult result;
try {
String nodeType = dag.getNodeTypeMap().get(currentNodeId);
NodeExecutor executor = executorMap.get(nodeType);
JsonNode nodeData = dag.getNodeDataMap().get(currentNodeId);
result = executor.execute(currentNodeId, nodeData, ctx, ws);
} catch (Exception e) {
result = NodeExecutionResult.failed(currentNodeId, "节点执行异常: " + e.getMessage());
}
// 4. 失败处理(沿用 DAG 模式的 failStrategy)
if (result.getStatus() == FAILED) {
String failStrategy = dag.getNodeDataMap().get(currentNodeId)
.path("failStrategy").asText("abort");
if (!"skip".equals(failStrategy)) {
// abort: 推 workflow_error,终止
abort(ctx, emitter, eventBus, runId, result.getError());
return;
}
// skip: 继续 Planner 决策
}
// 5. 合并输出到 currentVars(最近覆盖)
if (result.getOutput() != null) {
NodeOutputEnvelope env = NodeOutputEnvelope.fromObject(result.getOutput());
if (env != null && env.getData() != null) {
currentVars.putAll(flattenWithPrefix(env.getData(), currentNodeId));
}
}
// 6. 持久化 Visit + 推送 nodeResult SSE
String visitUuid = UUID.randomUUID().toString();
WorkflowRunVisit visit = buildVisit(
run.getId(), visitUuid, currentNodeId, revisits,
fromVisitUuid, fromNodeId, currentVars, result, sortOrder);
visitRepo.save(visit);
safeSend(emitter, WorkflowRunEvent.visitRecorded(runId, visit, result));
// 7. 终止判定:到达 Output 节点
String nodeType = dag.getNodeTypeMap().get(currentNodeId);
if ("output".equals(nodeType) && result.getStatus() == SUCCESS) {
ctx.setFinalOutput(result.getOutput());
run.setVisitCount(sortOrder + 1);
complete(ctx, emitter, eventBus, runId, result.getOutput());
return;
}
// 8. 调用 Planner 决定下一个节点
PlannerDecision decision;
try {
decision = plannerService.decideNext(
dag, currentNodeId, currentVars, nodeVisitCounter,
ctx.getInitialInputs(), aiModelService, fallbackBuilder);
} catch (Exception e) {
log.warn("[Agentic] Planner 失败,走画布参考路径: {}", e.getMessage());
decision = PlannerService.fallbackToNextEdge(dag, currentNodeId);
}
safeSend(emitter, WorkflowRunEvent.plannerDecision(runId, decision));
// 9. 校验 next 节点合法
if (!dag.getNodeDataMap().containsKey(decision.getNextNodeId())) {
log.warn("[Agentic] Planner 返回非法节点 {},走画布参考路径", decision.getNextNodeId());
decision = PlannerService.fallbackToNextEdge(dag, currentNodeId);
}
// 10. 准备下一次迭代
fromNodeId = currentNodeId;
fromVisitUuid = visitUuid;
currentNodeId = decision.getNextNodeId();
sortOrder++;
}
// 11. maxVisits 用尽
abort(ctx, emitter, eventBus, runId,
"超过最大访问次数 " + MAX_VISITS + ",疑似死循环");
}
/** 找起始节点:userInput 类型优先;否则取入度为 0 的第一个 */
private String findStartNode(ResolvedDag dag) { /* ... */ }
/** envelope.data 扁平化为 "{nodeId}__{varName}" 形式 */
private Map<String, Object> flattenWithPrefix(Map<String, Object> data, String nodeId) {
Map<String, Object> result = new LinkedHashMap<>();
data.forEach((k, v) -> result.put(nodeId + "__" + k, v));
return result;
}
}
新文件:backend/src/main/java/com/agent/management/engine/PlannerService.java
@Slf4j
@Service
@RequiredArgsConstructor
public class PlannerService {
private final ObjectMapper objectMapper;
public PlannerDecision decideNext(ResolvedDag dag, String currentNodeId,
Map<String, Object> currentVars,
Map<String, Integer> nodeVisitCounter,
Map<String, Object> initialInputs,
AiModelService aiModelService,
ChatClient.Builder fallbackBuilder) {
// 1. 构建 prompt
String prompt = buildPrompt(dag, currentNodeId, currentVars, nodeVisitCounter);
// 2. 调 LLM(最多重试 1 次)
for (int attempt = 0; attempt < 2; attempt++) {
try {
String response = aiModelService.getChatClientOrFallback(fallbackBuilder)
.prompt().user(prompt).call().content();
PlannerDecision decision = parseDecision(response, dag);
if (decision != null) return decision;
} catch (Exception e) {
log.warn("[Planner] 第 {} 次尝试失败: {}", attempt + 1, e.getMessage());
}
}
// 3. 全部失败,抛异常给上层走 fallback
throw new IllegalStateException("Planner 决策失败");
}
private String buildPrompt(ResolvedDag dag, String currentNodeId,
Map<String, Object> currentVars,
Map<String, Integer> nodeVisitCounter) {
// 节点清单
StringBuilder nodes = new StringBuilder();
dag.getNodeDataMap().forEach((id, data) -> {
String type = dag.getNodeTypeMap().get(id);
String label = data.path("label").asText(id);
int visited = nodeVisitCounter.getOrDefault(id, 0);
nodes.append("- ").append(id).append(" (").append(type).append("/").append(label).append(")");
if (visited > 0) nodes.append(" [已访问 ").append(visited).append(" 次]");
nodes.append("\n");
});
// 工作区变量预览(裁剪到 Top-30,每个值截断到 200 字符)
String varsPreview = currentVars.entrySet().stream()
.limit(30)
.map(e -> e.getKey() + ": " + truncate(String.valueOf(e.getValue()), 200))
.collect(Collectors.joining("\n"));
return """
你是工作流调度助手。请根据当前工作区状态决定下一步访问哪个节点。
当前刚执行完:%(currentNode)s
工作区变量(部分):
%(vars)s
可访问节点清单:
%(nodes)s
规则:
1. 选择 output 节点表示产出最终结果并结束运行
2. 不要在同一节点上反复横跳(已访问次数多的节点优先级降低)
3. 必须从上面的节点清单中选择
请回复 JSON:{"next": "节点id", "reason": "50字以内的理由"}
""".formatted(...);
}
private PlannerDecision parseDecision(String response, ResolvedDag dag) {
// 提取 JSON(容忍 LLM 加 ```json ``` 包裹)
String json = extractJson(response);
try {
JsonNode node = objectMapper.readTree(json);
String next = node.path("next").asText("");
String reason = node.path("reason").asText("");
if (next.isEmpty()) return null;
return new PlannerDecision(next, reason);
} catch (Exception e) {
return null;
}
}
/** 兜底:走画布上 currentNodeId 的第一条出边 */
public static PlannerDecision fallbackToNextEdge(ResolvedDag dag, String currentNodeId) {
List<DagResolver.EdgeInfo> edges = dag.getOutgoingEdges().get(currentNodeId);
if (edges != null && !edges.isEmpty()) {
return new PlannerDecision(edges.get(0).getTarget(),
"Planner 失败,走画布参考路径");
}
// 没有出边,找画布上任意一个 output 节点
for (Map.Entry<String, String> e : dag.getNodeTypeMap().entrySet()) {
if ("output".equals(e.getValue())) {
return new PlannerDecision(e.getKey(), "Planner 失败,强制结束到 output");
}
}
throw new IllegalStateException("无可达终止节点");
}
}
文件:backend/src/main/java/com/agent/management/engine/NodeWorkspaceBuilder.java
新增方法:
/**
* AGENTIC 模式专用:跳过可达前驱计算,直接用 currentVars 作为 workspace.variables。
* 输入解析(NodeInputResolver)仍正常工作,从这些变量里取值。
*/
public NodeWorkspace buildForAgentic(String nodeId, ResolvedDag dag,
WorkflowContext ctx,
Map<String, Object> currentVars) {
Map<String, Object> vars = new LinkedHashMap<>(currentVars);
return NodeWorkspace.builder()
.nodeId(nodeId)
.variables(vars)
.scopedOutputs(extractScopedFromVars(vars))
.build();
}
文件:backend/src/main/java/com/agent/management/engine/WorkflowEngine.java
execute() 入口加判断:
if ("AGENTIC".equalsIgnoreCase(workflow.getMode())) {
agenticExecutor.execute(dag, ctx, run, emitter, eventBus);
} else {
levelExecutor.executeByLevel(dag, ctx, run, emitter, eventBus);
}
文件:backend/src/main/java/com/agent/management/engine/WorkflowRunEvent.java
新增两个工厂方法:
public static WorkflowRunEvent visitRecorded(String runId, WorkflowRunVisit visit,
NodeExecutionResult result) {
Map<String, Object> data = new LinkedHashMap<>();
data.put("visitId", visit.getVisitUuid());
data.put("nodeId", visit.getNodeId());
data.put("nodeType", visit.getNodeType());
data.put("nodeLabel", visit.getNodeLabel());
data.put("iterSeq", visit.getIterSeq());
data.put("fromNodeId", visit.getFromNodeId());
data.put("fromVisitId", visit.getFromVisitUuid());
data.put("status", visit.getStatus());
data.put("output", result.getOutput());
data.put("sortOrder", visit.getSortOrder());
return new WorkflowRunEvent(runId, "visit_recorded", data, System.currentTimeMillis());
}
public static WorkflowRunEvent plannerDecision(String runId, PlannerDecision decision) {
Map<String, Object> data = new LinkedHashMap<>();
data.put("next", decision.getNextNodeId());
data.put("reason", decision.getReason());
return new WorkflowRunEvent(runId, "planner_decision", data, System.currentTimeMillis());
}
新文件:backend/src/main/resources/db/migration/V<下一个版本号>__add_agentic_mode.sql
-- 1. workflows 加 mode 字段
ALTER TABLE workflows ADD COLUMN mode VARCHAR(16) NOT NULL DEFAULT 'DAG';
-- 2. workflow_runs 加 visit_count
ALTER TABLE workflow_runs ADD COLUMN visit_count INT NULL;
-- 3. 新建 workflow_run_visits 表
CREATE TABLE workflow_run_visits (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
run_id BIGINT NOT NULL,
visit_uuid VARCHAR(64) NOT NULL UNIQUE,
node_id VARCHAR(128) NOT NULL,
node_type VARCHAR(64),
node_label VARCHAR(256),
iter_seq INT NOT NULL,
from_visit_uuid VARCHAR(64),
from_node_id VARCHAR(128),
variables_json TEXT,
output_json TEXT,
decision_json TEXT,
status VARCHAR(16) NOT NULL,
error_msg TEXT,
sort_order INT NOT NULL,
created_at DATETIME,
INDEX idx_run_id (run_id),
INDEX idx_run_node (run_id, node_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
注:若项目用 Flyway 自动按版本号执行;若用 JPA ddl-auto=update,可省略建表脚本。
文件:backend/src/main/java/com/agent/management/controller/WorkflowController.java
mode 字段GET /api/workflows/{id}/runs/{runId} 返回时附带 visits 列表(AGENTIC 模式)新增端点:
@GetMapping("/workflows/{workflowId}/runs/{runId}/visits")
public Result<List<VisitDto>> listVisits(@PathVariable Long runId) {
return Result.success(visitRepo.findByRunIdOrderBySortOrderAsc(runId)
.stream().map(VisitDto::from).toList());
}
新文件:backend/src/main/java/com/agent/management/dto/VisitDto.java
public record VisitDto(
String visitId,
String nodeId,
String nodeType,
String nodeLabel,
Integer iterSeq,
String fromVisitId,
String fromNodeId,
Map<String, Object> variables,
Map<String, Object> output,
Map<String, Object> decision,
String status,
String error,
Integer sortOrder,
String createdAt
) {
public static VisitDto from(WorkflowRunVisit v) { /* ... */ }
}
文件:frontend/src/components/workflow/BasicInfoDrawer.vue
<n-form-item label="运行模式">
<n-radio-group v-model:value="form.mode">
<n-radio value="DAG">DAG(按画布连线严格执行)</n-radio>
<n-radio value="AGENTIC">智能体(AI 自主决定跳转)</n-radio>
</n-radio-group>
</n-form-item>
切换到 AGENTIC 时弹确认:
function onModeChange(val) {
if (val === 'AGENTIC') {
dialog.info({
title: '智能体模式说明',
content: '此模式下,画布作为能力地图,AI 可访问任意节点。同节点可能被多次访问,运行结果会按访问顺序展开。',
positiveText: '我知道了'
})
}
}
文件:frontend/src/views/workflow/WorkflowEditor.vue
画布右上角加角标:
<div v-if="mode === 'AGENTIC'" class="agentic-badge">AGENTIC MODE</div>
文件:frontend/src/composables/useWorkflowRunner.js
const visits = ref([])
const plannerDecisions = ref([])
function handleSSEEvent(evt) {
switch (evt.type) {
case 'visit_recorded':
visits.value.push(evt.data)
// 同步更新 nodeStatuses / nodeResults(兼容现有 UI)
nodeStatuses.value[evt.data.nodeId] = evt.data.status
break
case 'planner_decision':
plannerDecisions.value.push(evt.data)
break
// ... 现有 case 不动
}
}
文件:
frontend/src/components/agent-thinking/AgentThinkingPanel.vue(编辑器右侧)frontend/src/views/workflow/RunHistory.vue(运行历史)判断当前 run 的 mode:
nodeRecords(不变)visits展示形态(AGENTIC 模式):
┌─ Visit 1:A (userInput) ──────────────────────┐
│ iter 1 · 起点 │
│ 工作区:{ user_query: "..." } │
│ 输出:{ user_query: "..." } │
└────────────────────────────────────────────────┘
↓ Planner 决策:转 B,因为需要 LLM 解析
┌─ Visit 2:B (llm) ────────────────────────────┐
│ iter 1 · ⬅ 来自 A#1 │
│ 工作区:{ A__user_query: "...", B__intent: "..." } │
│ 输出:{ intent: "查询天气" } │
└────────────────────────────────────────────────┘
↓ Planner 决策:转 C,因为需要技能调用
┌─ Visit 3:C (skill) ──────────────────────────┐
│ iter 1 · ⬅ 来自 B#1 │
│ 工作区:{ ..., C__result: "北京今天晴" } │
└────────────────────────────────────────────────┘
↓ Planner 决策:转 A,因为需要补充信息
┌─ Visit 4:A (userInput) ──────────────────────┐
│ iter 2 · ⬅ 来自 C#1 │
│ 工作区:{ ..., A__user_query: "那明天呢?" } │ ← A#1 同名字段被覆盖
└────────────────────────────────────────────────┘
↓ Planner 决策:转 output,任务完成
┌─ Visit 5:output ─────────────────────────────┐
│ ⬅ 来自 A#2 │
│ 最终输出:{ reply: "北京明天多云" } │
└────────────────────────────────────────────────┘
⬅ 来自 C#1新文件:frontend/src/components/workflow/PlannerDecisionTrail.vue
折叠面板,展示每次决策的"原因",便于事后追溯。
新文件:backend/src/test/java/com/agent/management/engine/
AgenticExecutorTest:
PlannerServiceTest:
VisitChainTest:
docs/plans/agentic-workflow-plan.md(即此文件)docs/design/agentic-workflow-design.md(实施完成后输出,给后续维护者看)| 风险 | 等级 | 缓解 |
|---|---|---|
| Planner LLM 决策失误,导致死循环 | 高 | ① maxVisits 硬上限 50;② 同节点连续访问 5 次强制终止;③ Planner prompt 中提示"已访问次数多的节点优先级降低" |
| Planner 输出 JSON 解析失败 | 中 | 重试 1 次,仍失败走 fallbackToNextEdge(画布第一条出边);都没有则找任意 output 节点 |
| 工作区变量无限膨胀 | 中 | Planner prompt 中只预览 Top-30;持久化时整体上限 64KB,超出截断标记 truncated |
| AGENTIC 模式性能(每步都调 LLM) | 中 | 工作流可配置 plannerModelId 用便宜模型;流式输出决策过程让用户感知进度 |
| visit 表数据爆炸 | 中 | 单 run 上限 50 visit;variables_json 单条上限 64KB |
| 现有 DAG 工作流兼容性 | 低 | mode 默认 DAG,新逻辑完全独立路径 |
| 节点输入解析(NodeInputResolver) | 低 | AGENTIC 模式下 workspace.variables 已被填满(带 {nodeId}__ 前缀),NodeInputResolver 的四级 fallback 照常工作 |
| Planner service 选择模型策略 | 低 | 复用 ConditionExecutor 的同款策略:aiModelService.getChatClientOrFallback(fallbackBuilder),行为与现有 ConditionNode 一致 |
| 阶段 | 内容 | 工作量 |
|---|---|---|
| 阶段 1 | 后端数据模型 | 1 人天 |
| 阶段 2 | 后端核心引擎 | 3 人天 |
| 阶段 3 | 数据库迁移 + API | 1 人天 |
| 阶段 4 | 前端编辑器 mode 切换 | 1 人天 |
| 阶段 5 | 前端运行结果展示 | 2 人天 |
| 阶段 6 | 测试 + 文档 | 2 人天 |
| 合计 | 10 人天 |
如果想要更小的 MVP 验证可行性,建议:
MVP 范围(5-6 天):
完整版(10 天):本文档全部内容
docs/plans/agentic-workflow-plan.md等待用户确认是否照此文档开始实施。
| 阶段 | 内容 | 状态 |
|---|---|---|
| 阶段 1 | 后端数据模型(Workflow.mode + WorkflowRunVisit) | ✅ |
| 阶段 2 | 后端 AgenticExecutor + PlannerService + SSE 事件扩展 | ✅ |
| 阶段 3 | Controller API(RunHistoryController 扩展,非新 Controller) | ✅ |
| 阶段 4 | 前端编辑器 mode 切换(基础信息 Tab 加 DAG / AGENTIC 单选) | ✅ |
| 阶段 5 | 前端运行结果 Visit 链展示(编辑器实时 + 运行记录回放) | ✅ |
| 阶段 6 | 后端单测(PlannerServiceTest + AgenticExecutorTest)+ 文档 | ✅ |
新增文件:
backend/src/main/java/com/agent/management/model/entity/WorkflowRunVisit.java — Visit 实体backend/src/main/java/com/agent/management/repository/WorkflowRunVisitRepository.java — JPA Repobackend/src/main/java/com/agent/management/engine/PlannerService.java — LLM 决策器(含 fallbackToNextEdge 静态兜底)backend/src/main/java/com/agent/management/engine/AgenticExecutor.java — AGENTIC 模式执行器backend/src/test/java/com/agent/management/engine/PlannerServiceTest.java — 11 个测试用例backend/src/test/java/com/agent/management/engine/AgenticExecutorTest.java — 6 个测试用例修改文件:
backend/.../model/entity/Workflow.java — 加 mode 字段backend/.../model/entity/WorkflowRun.java — 加 visitCount 字段backend/.../engine/WorkflowRunEvent.java — 加 visitRecorded / plannerDecision 工厂方法backend/.../engine/WorkflowEngine.java — mode 分流(AGENTIC 走 AgenticExecutor)backend/.../controller/RunHistoryController.java — RunDetailVO 扩展 mode / visits 字段backend/.../model/vo/WorkflowVO.java — 加 mode 字段(默认 "DAG")backend/.../service/WorkflowService.java (+Impl) — create/update 签名加 mode 参数backend/.../controller/WorkflowController.java — DTO 加 modefrontend/.../composables/useWorkflowRunner.js — visits / plannerDecisions 响应式状态 + SSE 事件分支frontend/.../views/workflow/WorkflowEditor.vue — 基础信息 Tab 加 mode 切换;运行结果 Tab 加 Visit 链frontend/.../views/workflow/RunHistory.vue — 加 Visit 链回放(按 mode 切换 DAG 节点列表 / AGENTIC 访问链)mvn test 全量 217 个测试,0 失败,5 跳过(与本次改造无关)
PlannerServiceTest:11/11 通过(fallback 4 个 + decideNext 6 个 + PlannerDecision 1 个)AgenticExecutorTest:6/6 通过(线性 / 工作区累积 / Planner 失败兜底 / 节点失败 abort / userInput 起始 / 死循环兜底)npm run build 通过,产物已在 backend/src/main/resources/static/| 计划方案 | 实施方案 | 原因 |
|---|---|---|
| 数据库迁移:可选 Flyway | 直接用现有 ddl-auto: update |
用户明确选择,避免引入额外依赖;Hibernate 自动建 workflow_run_visit 表与 workflow.mode / workflow_run.visit_count 列 |
| 新增 VisitController | 扩展现有 RunHistoryController |
单一 API 调用获取 nodes + visits,前端无需新增请求 |
| 前端 Visit 链抽公共组件 | 在 WorkflowEditor.vue 与 RunHistory.vue 各自内联 |
两处 visit 数据来源不同(实时 SSE vs 历史回放),抽象反而增加复杂度,遵循 KISS |
| VisitDto 数据转换 | 直接序列化 WorkflowRunVisit 实体 |
字段名已经合规(驼峰),无需 DTO 转换层 |
| 限制 | 值 | 触发后行为 |
|---|---|---|
| MAX_VISITS | 50 | abort,errorMsg 含"超过最大访问次数" |
| MAX_SAME_NODE_REVISITS | 5 | abort,errorMsg 含"死循环"。注:触发时第 6 次 visit 未保存,最终持久化 5 条 visit |
| Planner 失败重试 | 2 次(首次 + 1 次重试) | 全部失败后抛 IllegalStateException,由 AgenticExecutor 捕获并走 fallbackToNextEdge |