📷 [图片 token=Wy2CbrxSOoMSeLxMPYGc9co4nSd(未能下载,见飞书原文)]

前言

关键代码:app/agent/aiops/ 目录下的 state.pyplanner.pyexecutor.pyreplanner.py,以及 app/services/aiops_service.py

📷 [图片 token=BxHebByEioqLYQxFy82chBk2nnh(未能下载,见飞书原文)]

流程梳理

运维 Agent 的核心目标是 规划 → 执行 → 评估 → 调整。整体流程就是三个节点:

  1. Planner:拆解排查步骤,生成执行计划

  2. Executor:从计划中取出第一个步骤,调用工具执行

  3. Replanner:评估执行结果,决定继续、调整计划还是生成最终报告

三个节点通过 LangGraph StateGraph 串联,共享同一份 PlanExecuteState 状态对象在整个流程中传递。

📷 [图片 token=S0bFbHxhmoplyuxMjrJcPpXhn0d(未能下载,见飞书原文)]

实战

状态定义

整个 Plan-Execute-Replan 流程的数据流通过 PlanExecuteState 承载,字段设计非常简洁:

class PlanExecuteState(TypedDict):
    input: str                                          # 用户输入的任务描述
    plan: List[str]                                     # 待执行的步骤列表
    past_steps: Annotated[List[tuple], operator.add]   # 已执行的步骤历史(追加式更新)
    response: str                                       # 最终报告/响应

past_steps 使用 Annotated[List[tuple], operator.add] 声明,LangGraph 会将每次节点返回的 past_steps 自动追加到列表中,而不是覆盖,无需手动维护历史。

Planner 节点

Planner 负责制定执行计划,输出一个结构化的步骤列表。核心流程:

  1. 先调用 retrieve_knowledge 查询知识库,寻找历史经验文档

  2. 获取所有可用工具(本地工具 + MCP 工具),格式化为文字描述

  3. 将工具列表和经验文档注入 prompt,调用 LLM 生成结构化计划

async def planner(state: PlanExecuteState) -> Dict[str, Any]:
    input_text = state.get("input", "")

    # 1. 查询内部知识库,寻找相关经验
    context_str = await retrieve_knowledge.ainvoke({"query": input_text})
    experience_docs = context_str if context_str and context_str.strip() else ""

    # 2. 获取所有可用工具(本地 + MCP)
    local_tools = [get_current_time, retrieve_knowledge]
    mcp_tools = await (await get_mcp_client_with_retry()).get_tools()
    all_tools = local_tools + mcp_tools

    # 3. 调用 LLM 生成结构化计划
    llm = ChatQwen(model=config.rag_model, api_key=config.dashscope_api_key, temperature=0)
    planner_chain = planner_prompt | llm.with_structured_output(Plan)

    plan_result = await planner_chain.ainvoke({
        "messages": [("user", input_text)],
        "tools_description": format_tools_description(all_tools),
        "experience_context": experience_docs
    })

    return {"plan": plan_result.steps}

计划的输出格式用 Pydantic Plan 模型约束,通过 llm.with_structured_output(Plan) 保证 LLM 的输出可以直接解析为步骤列表:

class Plan(BaseModel):
    steps: List[str] = Field(
        description="完成任务所需的不同步骤,按顺序执行,每一步建立在前一步的基础上。"
    )

Planner Prompt

Planner 的系统提示词要求模型将任务分解为逻辑独立的步骤,每步指明使用哪个工具及所需参数。如果查到了经验文档,也会作为参考注入:

planner_prompt = ChatPromptTemplate.from_messages([
    ("system", """
        作为一个专家级别的规划者,你需要将复杂的任务分解为可执行的步骤。

        可用工具列表(用于制定计划时参考):
        {tools_description}

        注意:你的职责是制定计划,实际的工具调用由 Executor 负责执行。

        {experience_context}

        对于给定的任务,请创建一个简单的、逐步的计划:
        - 将任务分解为逻辑上独立的步骤
        - 每个步骤明确使用哪些工具(如果需要),最好能同时提供工具所需参数
        - 步骤之间应有清晰的依赖关系
        - 如果有相关经验文档,请参考其中的方法和步骤制定计划
    """),
    ("placeholder", "{messages}"),
])

Executor 节点

Executor 每次只执行计划中的第一个步骤,使用 LangGraph 的 ToolNode 自动处理工具调用,执行完后将该步骤从 plan 中移除,并将执行结果追加到 past_steps

async def executor(state: PlanExecuteState) -> Dict[str, Any]:
    plan = state.get("plan", [])
    task = plan[0]   # 只取第一个步骤

    # 绑定工具的 LLM
    all_tools = local_tools + mcp_tools
    llm_with_tools = llm.bind_tools(all_tools)
    tool_node = ToolNode(all_tools)

    messages = [
        SystemMessage(content="你是一个能力强大的助手,负责执行具体的任务步骤。..."),
        HumanMessage(content=f"请执行以下任务: {task}")
    ]

    # 第一步:LLM 决定是否需要工具调用
    llm_response = await llm_with_tools.ainvoke(messages)

    # 第二步:如果有工具调用,使用 ToolNode 自动执行
    if hasattr(llm_response, "tool_calls") and llm_response.tool_calls:
        messages.append(llm_response)
        tool_messages = await tool_node.ainvoke({"messages": messages})
        messages.extend(tool_messages["messages"])
        # 第三步:将工具结果返回给 LLM 生成最终答案
        final_response = await llm_with_tools.ainvoke(messages)
        result = final_response.content
    else:
        result = llm_response.content

    return {
        "plan": plan[1:],               # 移除已执行的第一个步骤
        "past_steps": [(task, result)], # 追加执行历史
    }

Replanner 节点

Replanner 根据原始任务、已执行步骤和剩余计划做出三选一的决策:

决策含义触发条件
respond信息充足,立即生成最终报告最高优先级,已执行 ≥ 3 步且有关键信息
continue当前计划合理,继续执行剩余步骤确实必要
replan调整计划,替换剩余步骤最低优先级,计划有重大偏差时才使用

决策同样用 Pydantic 模型约束输出:

class Act(BaseModel):
    action: str = Field(description="下一步行动: 'continue' | 'replan' | 'respond'")
    new_steps: List[str] = Field(default_factory=list, description="replan 时的新步骤列表")

Replanner 核心逻辑(含安全限制):

async def replanner(state: PlanExecuteState) -> Dict[str, Any]:
    past_steps = state.get("past_steps", [])
    plan = state.get("plan", [])

    # 安全限制:已执行步骤超过 8 步,强制生成响应,防止无限循环
    if len(past_steps) >= 8:
        return await _generate_response(state, llm)

    if plan:
        # 还有剩余步骤,让 LLM 做决策
        act = await replanner_chain.ainvoke({
            "messages": [
                ("user", f"原始任务: {input_text}"),
                ("user", f"已执行的步骤:\n{steps_summary}"),
                ("user", f"剩余计划: {', '.join(plan)}"),
            ],
            "tools_description": tools_description
        })

        if act.action == "respond":
            return await _generate_response(state, llm)
        elif act.action == "replan":
            # 安全限制:新步骤数不能超过当前剩余步骤数
            new_steps = act.new_steps[:len(plan)]
            return {"plan": new_steps}
        else:  # continue
            return {}   # 不修改状态,继续执行下一步
    else:
        # 计划已执行完毕,直接生成最终响应
        return await _generate_response(state, llm)

当决定 respond 时,_generate_response 会整理所有执行历史,生成结构化的 Markdown 报告:

async def _generate_response(state, llm) -> Dict[str, Any]:
    execution_history = "\n\n".join([
        f"### 步骤: {step}\n**结果:**\n{result}"
        for step, result in past_steps
    ])

    response_gen = response_prompt | llm.with_structured_output(Response)
    response_obj = await response_gen.ainvoke({
        "messages": [
            ("user", f"原始任务: {input_text}"),
            ("user", f"执行历史:\n{execution_history}"),
            ("user", "请基于以上信息生成全面的最终响应"),
        ]
    })
    return {"response": response_obj.response}

构建 LangGraph 工作流

三个节点通过 StateGraph 连接,replanner 之后根据状态中是否存在 response 进行条件路由:

def _build_graph(self):
    workflow = StateGraph(PlanExecuteState)

    # 添加三个节点
    workflow.add_node("planner", planner)
    workflow.add_node("executor", executor)
    workflow.add_node("replanner", replanner)

    # 固定边:planner -> executor -> replanner
    workflow.set_entry_point("planner")
    workflow.add_edge("planner", "executor")
    workflow.add_edge("executor", "replanner")

    # 条件边:replanner 根据状态决定结束还是继续执行
    def should_continue(state: PlanExecuteState) -> str:
        if state.get("response"):
            return END          # 已生成最终响应,结束流程
        if state.get("plan"):
            return "executor"   # 还有步骤,继续执行
        return END

    workflow.add_conditional_edges("replanner", should_continue, {
        "executor": "executor",
        END: END
    })

    return workflow.compile(checkpointer=MemorySaver())

工作流结构如下:

[用户任务] → planner → executor → replanner
                            ↑           |
                            |   continue/replan
                            └───────────┘
                                        |
                                   respond/空计划
                                        ↓
                                      [END]

执行工作流(流式输出)

AIOpsService.execute 使用 graph.astream(stream_mode="updates") 流式执行,每个节点完成后立即产生事件,可以实时推送给前端展示执行进度:

async def execute(self, user_input: str, session_id: str) -> AsyncGenerator:
    initial_state: PlanExecuteState = {
        "input": user_input,
        "plan": [],
        "past_steps": [],
        "response": ""
    }

    async for event in self.graph.astream(
        input=initial_state,
        config={"configurable": {"thread_id": session_id}},
        stream_mode="updates"   # 每个节点更新后立即推送
    ):
        for node_name, node_output in event.items():
            if node_name == "planner":
                yield self._format_planner_event(node_output)     # type: plan
            elif node_name == "executor":
                yield self._format_executor_event(node_output)    # type: step_complete
            elif node_name == "replanner":
                yield self._format_replanner_event(node_output)   # type: report / status

    # 所有节点完成,推送最终响应
    final_state = self.graph.get_state({"configurable": {"thread_id": session_id}})
    final_response = final_state.values.get("response", "") if final_state else ""

    yield {
        "type": "complete",
        "stage": "complete",
        "message": "任务执行完成",
        "response": final_response
    }

流式事件类型说明:

typestage含义
planplan_createdPlanner 生成了执行计划
step_completestep_executedExecutor 执行完一个步骤
reportfinal_reportReplanner 生成了最终报告
status各节点名节点运行中的状态通知
completecomplete整个工作流结束
errorerror执行出错

总结

通过上面的分析,我们已经了解了 Planner、Executor、Replanner 的作用和相关 prompt。代码的核心是 LangGraph 的 StateGraph 管理节点间的状态流转,PlanExecuteState 在整个流程中传递,三个节点各司其职:

  1. Planner 查询知识库获取经验,生成结构化步骤列表

  2. Executor 取出第一步,通过 ToolNode 自动执行工具调用,返回结果并移除该步

  3. Replanner 评估已执行结果,决定是继续、调整还是收敛生成最终报告

框架帮我们处理了节点间的数据传递、条件路由和流式输出,核心还是要搞懂设计原理:Plan-Execute-Replan 本质上是一个带反馈的闭环调度器,Replanner 的收敛策略决定了整个流程的质量与效率。

📷 [图片 token=GIrDbzR7WouUWMxHlD3cQGyYnpf(未能下载,见飞书原文)]