agentic-workflow-plan.md 37 KB

自主智能体工作流模式(Agentic Mode)实施计划

文档目的:在现有 DAG 工作流引擎基础上,新增一种"自主智能体"执行模式。在该模式下,画布作为智能体的"能力地图"(参考路径,非约束),由全局 Planner Agent 决定每一步访问哪个节点;同一节点可被多次访问,每次访问独立暂存工作区快照。

改造范围

  • 后端:新增 AgenticExecutor + PlannerService + WorkflowRunVisit 实体(独立路径,不动 DAG 引擎)
  • 前端:编辑器加 mode 切换、运行结果改为按 Visit 链展示

状态:待实施


1. 问题确认

1.1 用户需求

我做的不是传统的工作流智能体,而是灵活的可自主决定如何运行的智能体,工作流只是给它的参考而非桎梏。所以,可以创造画布上不存在的边,只要决策 AI 觉得有必要。对应的,每个节点,由于可能经过多次,所以要暂存每次运行时的工作区;若下一次运行(从后继节点飞回来),则将后继节点的工作区变量迁移过来,再次记录。相应地,在"运行结果"处,也会多次出现该节点,每次出现时,工作区变量都可能不同。 如果出现灵活运行的边,要在运行结果处给飞到的节点打一个标识,说明它是从哪个节点飞过来的。

1.2 范式对比

维度 DAG 模式(现有) AGENTIC 模式(本次新增)
画布连线 必须严格遵守的控制流 参考路径,非约束
节点语义 固定步骤,执行一次 能力/工具,可被任意次访问
引擎角色 调度器("该你了") 执行体(决策者另有其人)
决策权 静态画布连线 全局 Planner Agent(运行时 LLM)
节点输出 单次写入命名空间 每次 Visit 独立快照
终止条件 所有节点完成 访问到 Output 节点

1.3 关键决策点(已与用户对齐)

决策 选择 原文
决策者 X2 全局 Planner "工作流只是给它的参考" 暗示有"它"主体
上下文迁移 链式继承 + 最近覆盖 用户给出明确语义:A#2 工作区 = C#1 工作区 + A#2 输出(覆盖 A#1 同名变量)
跳转约束 画布节点集合内任意 W2,平衡灵活性和可调试性
终止条件 必须到 Output T1,保证有明确产出

1.4 Visit 链式记忆语义(按用户原话)

访问序列: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 工作区 + 本节点输出,同名变量最近覆盖。

1.5 仓库现状

  • 现有 WorkflowLevelExecutor 严格执行 DAG,禁止环(DagResolver.java:108-110 检测到环抛异常)
  • 已有 ConditionExecutor.java 是 LLM 驱动的分支选择(仅在画布已画的出边里选),可作为 Planner prompt 的实现参考
  • 节点输出统一封装为 NodeOutputEnvelope,节点级命名空间隔离在 WorkflowContext.nodeScopedOutputs
  • SSE 推送走 SseEventBus(双轨:emitter + 缓存)支持断线重连

本次为全新独立路径,不动 DAG 引擎任何代码。


2. 改造目标

2.1 行为目标

编号 行为
B1 工作流可配置 mode = DAG(默认)/ AGENTIC
B2 AGENTIC 模式下,运行时由全局 Planner Agent 决定下一步访问哪个节点
B3 Planner 可指定画布上任意节点(不要求画布有连线),也可指定 Output 节点终结
B4 每个节点访问生成独立的 Visit 快照,含工作区变量、输出、Planner 决策原因
B5 运行结果按 Visit 顺序展开,同节点多次访问逐条显示,飞回的 Visit 标识来源
B6 兜底:maxVisits 硬上限 + Planner 失败走画布参考路径

2.2 必须保留的语义

编号 现有行为 保留方式
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

3. 设计选型

方案 A:独立 AgenticExecutor(已选

新建 AgenticExecutor 类,与 WorkflowLevelExecutor 并列,在 WorkflowEngine.execute() 入口按 mode 分流。

优点

  • DAG 引擎零侵入,现有工作流零影响
  • AGENTIC 模式独立演进,不影响 DAG
  • 失败回退容易:mode 切回 DAG 即可

缺点

  • 部分逻辑(如 SSE 推送、NodeWorkspace 构建)有重复
  • 接受重复,换取隔离性

方案 B:在 WorkflowLevelExecutor 加分支(未选)

破坏单一职责,调试时难以分辨"这段代码是 DAG 还是 AGENTIC 在跑"。

方案 C:在节点层面循环(未选)

每节点自己决定下一步,违反"全局决策者"语义。


4. 实施阶段

阶段 1:后端 - 数据模型与持久化(~1 天)

1.1 Workflow 实体加 mode 字段

文件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";

1.2 新增 WorkflowRunVisit 实体

新文件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;
}

1.3 新增 Repository

新文件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);
}

1.4 WorkflowRun 加 visitCount

文件backend/src/main/java/com/agent/management/model/entity/WorkflowRun.java

@Column(name = "visit_count")
private Integer visitCount;  // AGENTIC 模式下的总访问次数;DAG 模式为 null

阶段 2:后端 - AgenticExecutor 核心引擎(~3 天)

2.1 新增 AgenticExecutor

新文件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;
    }
}

2.2 新增 PlannerService

新文件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("无可达终止节点");
    }
}

2.3 NodeWorkspaceBuilder 增加 AGENTIC 入口

文件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();
}

2.4 WorkflowEngine 分流

文件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);
}

2.5 SSE 事件扩展

文件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());
}

阶段 3:数据库迁移 + Controller API(~1 天)

3.1 数据库迁移脚本

新文件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,可省略建表脚本。

3.2 Controller API 扩展

文件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());
}

3.3 VisitDto

新文件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) { /* ... */ }
}

阶段 4:前端 - 编辑器加 mode 切换(~1 天)

4.1 BasicInfoDrawer 加运行模式单选

文件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: '我知道了'
    })
  }
}

4.2 编辑器画布显示 AGENTIC 徽标

文件frontend/src/views/workflow/WorkflowEditor.vue

画布右上角加角标:

<div v-if="mode === 'AGENTIC'" class="agentic-badge">AGENTIC MODE</div>

阶段 5:前端 - 运行结果 Visit 链展示(~2 天)

5.1 useWorkflowRunner 处理新事件

文件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 不动
  }
}

5.2 运行展示组件改造

文件

  • frontend/src/components/agent-thinking/AgentThinkingPanel.vue(编辑器右侧)
  • frontend/src/views/workflow/RunHistory.vue(运行历史)

判断当前 run 的 mode:

  • DAG 模式:数据源用 nodeRecords(不变)
  • AGENTIC 模式:数据源用 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: "北京明天多云" }             │
└────────────────────────────────────────────────┘

5.3 飞回标识样式

  • 标题左侧加蓝色徽标 ⬅ 来自 C#1
  • 同节点多次访问,iter 角标用不同颜色(A#1 灰、A#2 橙、A#3 红,警示作用)

5.4 Planner 决策链折叠面板(可选增强)

新文件frontend/src/components/workflow/PlannerDecisionTrail.vue

折叠面板,展示每次决策的"原因",便于事后追溯。


阶段 6:测试与文档(~2 天)

6.1 后端单元测试

新文件backend/src/test/java/com/agent/management/engine/

  • AgenticExecutorTest
    • 正常路径(A→B→output)
    • 飞回(A→B→A→output)
    • maxVisits 兜底
    • 同节点 5 次兜底
    • Planner 失败走画布参考路径
  • PlannerServiceTest
    • JSON 解析(含 ```json 包裹)
    • 非法节点名兜底
    • LLM 异常重试
  • VisitChainTest
    • 工作区链式继承
    • 同名变量最近覆盖

6.2 集成测试

  • 跑一个完整 AGENTIC 模式工作流
  • 断言 visits 表记录数、iterSeq 自增、飞回标识、finalOutput 正确

6.3 文档

  • 本计划文档docs/plans/agentic-workflow-plan.md(即此文件)
  • 设计文档docs/design/agentic-workflow-design.md(实施完成后输出,给后续维护者看)

5. 风险评估

风险 等级 缓解
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 一致

6. 复杂度估算

阶段 内容 工作量
阶段 1 后端数据模型 1 人天
阶段 2 后端核心引擎 3 人天
阶段 3 数据库迁移 + API 1 人天
阶段 4 前端编辑器 mode 切换 1 人天
阶段 5 前端运行结果展示 2 人天
阶段 6 测试 + 文档 2 人天
合计 10 人天

7. 关键不变式(实施时严格遵守)

  1. mode = DAG 行为完全不变 —— 所有现有工作流零影响,所有现有测试零改动通过
  2. Planner 决策失败必有兜底 —— 绝不抛异常中断运行,至少走画布参考路径或找 output 终结
  3. maxVisits 硬上限 50 —— 任何情况都强制终止,避免死循环
  4. 同节点连续访问 5 次强制终止 —— 防止 Planner 卡在两个节点间反复横跳
  5. Visit 链式记忆严格按"上一 visit 工作区 + 本节点输出" —— 不允许跨 visit 偷渡变量
  6. Output 节点是唯一合法终点 —— 其他节点都不能直接产出 finalOutput
  7. 同节点多次访问必须可见 —— 运行结果不合并、不折叠,逐条展示
  8. Planner 决策必须可追溯 —— 每个 visit 持久化 decisionJson 含 next + reason

8. MVP 优先实施顺序

如果想要更小的 MVP 验证可行性,建议:

MVP 范围(5-6 天)

  • 阶段 1 数据模型(核心字段)
  • 阶段 2 核心引擎(不含 maxSameNodeRevisits 兜底)
  • 阶段 3 数据库迁移 + 基础 API
  • 阶段 4 基础切换
  • 阶段 5.1 + 5.2 基础展示
  • 阶段 6 后端核心单测

完整版(10 天):本文档全部内容


9. 实施前的最终检查清单

  • 计划已写入 docs/plans/agentic-workflow-plan.md
  • 用户已确认需求、决策点、风险
  • 现有 DAG 引擎行为零影响(mode 默认 DAG)
  • Planner 兜底策略明确(fallbackToNextEdge)
  • maxVisits / maxSameNodeRevisits 上限值已定(50 / 5)
  • 数据库迁移方案已定(按项目是否用 Flyway 决定)

等待用户确认是否照此文档开始实施。


10. 实施记录(2026-08-08 完成)

10.1 实施概览

阶段 内容 状态
阶段 1 后端数据模型(Workflow.mode + WorkflowRunVisit)
阶段 2 后端 AgenticExecutor + PlannerService + SSE 事件扩展
阶段 3 Controller API(RunHistoryController 扩展,非新 Controller)
阶段 4 前端编辑器 mode 切换(基础信息 Tab 加 DAG / AGENTIC 单选)
阶段 5 前端运行结果 Visit 链展示(编辑器实时 + 运行记录回放)
阶段 6 后端单测(PlannerServiceTest + AgenticExecutorTest)+ 文档

10.2 关键文件清单

新增文件

  • backend/src/main/java/com/agent/management/model/entity/WorkflowRunVisit.java — Visit 实体
  • backend/src/main/java/com/agent/management/repository/WorkflowRunVisitRepository.java — JPA Repo
  • backend/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.javaRunDetailVO 扩展 mode / visits 字段
  • backend/.../model/vo/WorkflowVO.java — 加 mode 字段(默认 "DAG")
  • backend/.../service/WorkflowService.java (+Impl) — create/update 签名加 mode 参数
  • backend/.../controller/WorkflowController.java — DTO 加 mode
  • frontend/.../composables/useWorkflowRunner.js — visits / plannerDecisions 响应式状态 + SSE 事件分支
  • frontend/.../views/workflow/WorkflowEditor.vue — 基础信息 Tab 加 mode 切换;运行结果 Tab 加 Visit 链
  • frontend/.../views/workflow/RunHistory.vue — 加 Visit 链回放(按 mode 切换 DAG 节点列表 / AGENTIC 访问链)

10.3 测试结果

  • 后端: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/

10.4 与计划的偏差

计划方案 实施方案 原因
数据库迁移:可选 Flyway 直接用现有 ddl-auto: update 用户明确选择,避免引入额外依赖;Hibernate 自动建 workflow_run_visit 表与 workflow.mode / workflow_run.visit_count
新增 VisitController 扩展现有 RunHistoryController 单一 API 调用获取 nodes + visits,前端无需新增请求
前端 Visit 链抽公共组件 WorkflowEditor.vueRunHistory.vue 各自内联 两处 visit 数据来源不同(实时 SSE vs 历史回放),抽象反而增加复杂度,遵循 KISS
VisitDto 数据转换 直接序列化 WorkflowRunVisit 实体 字段名已经合规(驼峰),无需 DTO 转换层

10.5 运行时硬上限(与 §7 一致)

限制 触发后行为
MAX_VISITS 50 abort,errorMsg 含"超过最大访问次数"
MAX_SAME_NODE_REVISITS 5 abort,errorMsg 含"死循环"。注:触发时第 6 次 visit 未保存,最终持久化 5 条 visit
Planner 失败重试 2 次(首次 + 1 次重试) 全部失败后抛 IllegalStateException,由 AgenticExecutor 捕获并走 fallbackToNextEdge

10.6 后续可能迭代点(不在本次范围内)

  1. Visit 链可视化:当前为线性列表,未来可在画布上叠加 visit 跳转动画
  2. Planner prompt 优化:当前 prompt 较朴素,后续可加入"工具调用历史"维度让决策更准
  3. Visit 工作区 diff 视图:展示"本次访问新增 / 覆盖了哪些变量",便于调试
  4. MAX_SAME_NODE_REVISITS 阈值:当前固定 5,未来可按节点类型差异化(llm 容忍高,skill 容忍低)