摘要:Agent Loop(智能体循环)是当前 LLM Agent 最核心的执行范式,而 Graph Engineer(图工程)则代表下一代智能体编排方向——用有向图显式建模状态、节点、边与条件分支。本文从专业视角系统讲解两者的原理、边界与融合方式,并给出一套可运行的 Python 实现:包含 ReAct 循环、反思循环、工具调用、状态机、DAG 编排、条件边、并行节点、检查点与恢复、人机协作、子图嵌套、可观测性与测试。
最早的 LLM 应用是单轮问答:
用户输入 -> LLM -> 输出但它有三个根本问题:
ReAct(Reasoning + Acting)范式解决了这个问题:
Thought -> Action -> Observation -> Thought -> ... -> Final Answer这个循环就是 Agent Loop。
Agent Loop 的价值:
能力 | 单轮问答 | Agent Loop |
|---|---|---|
多步任务 | ❌ | ✅ |
工具调用 | ❌ | ✅ |
自我修正 | ❌ | ✅ |
长任务 | ❌ | ✅ |
可观测 | 弱 | 强 |
一个最小 Agent Loop:
while not done:
thought = llm.think(state)
if thought.is_final:
break
action = thought.action
observation = tools.run(action)
state.append(observation)Agent Loop 简单,但有明显上限:
Graph Engineer 的思路是:把 Agent 执行过程显式建模为图。
节点(Node)= 一个执行单元(LLM 调用、工具调用、函数)
边(Edge) = 控制流(顺序、条件、并行)
状态(State)= 在节点间流动的共享数据
检查点(Checkpoint)= 每一步的状态快照Graph 的优势:
能力 | Agent Loop | Graph Engineer |
|---|---|---|
状态可见 | 弱 | 强 |
并行 | ❌ | ✅ |
条件分支 | 隐式 | 显式 |
断点恢复 | ❌ | ✅ |
人机协作 | 难 | 原生支持 |
多智能体 | 难 | 子图 |
可视化 | 难 | 天然 |
一句话:Loop 是运行时,Graph 是编排层。
两者不是替代关系,而是互补:
┌─────────────────────────────────────────┐
│ Graph 编排层 │
│ ┌──────────┐ ┌──────────┐ ┌────────┐ │
│ │ Node A │->│ Node B │->│ Node C │ │
│ │ (Loop) │ │ (Tool) │ │ (LLM) │ │
│ └──────────┘ └──────────┘ └────────┘ │
└─────────────────────────────────────────┘工程上常见组合:
┌─────────────────────────────────────────┐
│ │
│ ┌───────────┐ │
│ │ Observe │◄────────────┐ │
│ └─────┬─────┘ │ │
│ ▼ │ │
│ ┌───────────┐ │ │
│ │ Think │ │ │
│ └─────┬─────┘ │ │
│ ▼ │ │
│ ┌───────────┐ │ │
│ │ Act │ │ │
│ └─────┬─────┘ │ │
│ ▼ │ │
│ ┌───────────┐ │ │
│ │ Check │───── 未完成 ┘ │
│ └─────┬─────┘ │
│ ▼ 完成 │
│ Final Answer │
└─────────────────────────────────────────┘变体 | 说明 | 适用场景 |
|---|---|---|
ReAct | Thought + Action + Observation | 通用工具调用 |
Reflexion | 增加自我反思节点 | 需要纠错 |
Plan-and-Execute | 先规划再执行 | 长任务 |
Tree-of-Thought | 多路径探索 | 复杂推理 |
Self-Consistency | 多次采样投票 | 提高稳定性 |
State(状态):节点间共享的数据
Node(节点):执行单元
Edge(边):控制流
Condition(条件):决定走哪条边
Checkpoint(检查点):状态快照1. 初始化 State
2. 从 START 节点开始
3. 执行当前节点,更新 State
4. 根据边与条件选择下一个节点
5. 重复直到到达 END
6. 每一步可保存检查点类型 | 说明 |
|---|---|
LLM Node | 调用大模型 |
Tool Node | 调用工具 |
Function Node | 纯函数处理 |
Condition Node | 条件判断 |
Parallel Node | 并行分支 |
Subgraph Node | 嵌套子图 |
Human Node | 等待人工输入 |
类型 | 说明 |
|---|---|
普通边 | 固定跳转 |
条件边 | 根据 State 选择 |
并行边 | 同时触发多个节点 |
汇聚边 | 等待多个节点完成 |
状态是一个可序列化的字典,包含:
每个节点是一个函数:
def node(state: State) -> State:
...
return updated_stategraph.add_edge("node_a", "node_b")
graph.add_conditional_edge("node_b", router_fn, {"yes": "node_c", "no": "node_d"})每步保存 State 快照,用于:
agent-graph-engine/
├── app/
│ ├── __init__.py
│ ├── main.py
│ ├── state.py
│ ├── loop.py
│ ├── tools.py
│ ├── graph.py
│ ├── nodes.py
│ ├── checkpoint.py
│ ├── hitl.py
│ ├── subgraph.py
│ ├── observability.py
│ └── llm.py
├── tests/
│ └── test_agent_graph.py
├── requirements.txt
└── Dockerfilepip install fastapi uvicorn pydantic sqlmodel pytest httpxrequirements.txt:
fastapi
uvicorn[standard]
pydantic
sqlmodel
pytest
httpxapp/state.py:
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional
from datetime import datetime
import uuid
@dataclass
class Message:
role: str # system / user / assistant / tool
content: str
tool_name: Optional[str] = None
tool_call_id: Optional[str] = None
timestamp: str = field(default_factory=lambda: datetime.utcnow().isoformat())
@dataclass
class AgentState:
"""
图执行的共享状态。所有节点读写同一个 State。
"""
run_id: str = field(default_factory=lambda: str(uuid.uuid4()))
user_id: int = 0
goal: str = ""
messages: List[Message] = field(default_factory=list)
variables: Dict[str, Any] = field(default_factory=dict)
tool_results: List[Dict[str, Any]] = field(default_factory=list)
errors: List[str] = field(default_factory=list)
step: int = 0
max_steps: int = 20
status: str = "running" # running / done / failed / waiting_human
next_node: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
def add_message(self, role: str, content: str, **kwargs):
self.messages.append(Message(role=role, content=content, **kwargs))
def set(self, key: str, value: Any):
self.variables[key] = value
def get(self, key: str, default: Any = None) -> Any:
return self.variables.get(key, default)
def to_dict(self) -> dict:
return {
"run_id": self.run_id,
"user_id": self.user_id,
"goal": self.goal,
"messages": [m.__dict__ for m in self.messages],
"variables": self.variables,
"tool_results": self.tool_results,
"errors": self.errors,
"step": self.step,
"status": self.status,
"next_node": self.next_node,
"metadata": self.metadata,
}
@classmethod
def from_dict(cls, data: dict) -> "AgentState":
state = cls(
run_id=data["run_id"],
user_id=data.get("user_id", 0),
goal=data.get("goal", ""),
variables=data.get("variables", {}),
tool_results=data.get("tool_results", []),
errors=data.get("errors", []),
step=data.get("step", 0),
status=data.get("status", "running"),
next_node=data.get("next_node"),
metadata=data.get("metadata", {}),
)
state.messages = [Message(**m) for m in data.get("messages", [])]
return stateapp/llm.py:
from dataclasses import dataclass
from typing import Optional
@dataclass
class LLMDecision:
kind: str # "tool" / "final" / "think"
content: str = ""
tool_name: Optional[str] = None
arguments: Optional[dict] = None
def mock_llm(
state,
tool_schemas: list,
) -> LLMDecision:
"""
Mock LLM:规则驱动,便于离线运行。
生产环境替换为 OpenAI / Claude function calling。
"""
goal = state.goal.lower()
last_msgs = state.messages[-5:]
# 看是否已经有工具结果,避免重复调用
has_tool_result = any(m.role == "tool" for m in state.messages)
if "计算" in state.goal and not has_tool_result:
return LLMDecision(
kind="tool",
tool_name="calculator",
arguments={"expression": "2+3*4"},
)
if "天气" in state.goal and not has_tool_result:
return LLMDecision(
kind="tool",
tool_name="weather",
arguments={"city": "上海"},
)
if "搜索" in state.goal and not has_tool_result:
return LLMDecision(
kind="tool",
tool_name="search",
arguments={"query": state.goal},
)
if has_tool_result:
return LLMDecision(
kind="final",
content=f"任务完成。基于工具结果,答案是:{state.tool_results[-1]}",
)
return LLMDecision(
kind="final",
content=f"收到目标:{state.goal}。当前没有合适工具,直接回答。",
)app/loop.py:
from app.state import AgentState
from app.llm import mock_llm
from app.tools import ToolRegistry
class AgentLoop:
"""
单节点内的 ReAct 循环。
"""
def __init__(
self,
tools: ToolRegistry,
max_steps: int = 10,
max_repeats: int = 2,
):
self.tools = tools
self.max_steps = max_steps
self.max_repeats = max_repeats
def run(self, state: AgentState) -> AgentState:
action_history = []
for step in range(self.max_steps):
state.step += 1
if state.step > state.max_steps:
state.status = "failed"
state.errors.append("超过最大步数")
return state
decision = mock_llm(state, self.tools.to_schema())
if decision.kind == "final":
state.add_message("assistant", decision.content)
state.status = "done"
return state
if decision.kind == "tool":
# 重复检测
key = f"{decision.tool_name}:{decision.arguments}"
action_history.append(key)
if action_history.count(key) > self.max_repeats:
state.errors.append(f"重复动作:{key}")
state.add_message("assistant", "检测到重复动作,停止循环。")
state.status = "failed"
return state
state.add_message(
"assistant",
f"调用工具 {decision.tool_name}",
tool_name=decision.tool_name,
)
result = self.tools.execute(
decision.tool_name,
decision.arguments or {},
state,
)
state.tool_results.append(result)
state.add_message(
"tool",
str(result),
tool_name=decision.tool_name,
)
continue
state.status = "failed"
state.errors.append("循环未收敛")
return stateapp/tools.py:
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, List
@dataclass
class ToolSpec:
name: str
description: str
parameters: Dict[str, Any]
handler: Callable[..., Any]
risk_level: str = "low"
class ToolRegistry:
def __init__(self):
self._tools: Dict[str, ToolSpec] = {}
def register(self, spec: ToolSpec):
self._tools[spec.name] = spec
def get(self, name: str) -> ToolSpec | None:
return self._tools.get(name)
def to_schema(self) -> List[dict]:
return [
{
"name": t.name,
"description": t.description,
"parameters": t.parameters,
}
for t in self._tools.values()
]
def execute(self, name: str, args: dict, state) -> dict:
spec = self.get(name)
if not spec:
return {"error": f"未知工具:{name}"}
try:
result = spec.handler(**args)
return {"tool": name, "ok": True, "result": result}
except Exception as e:
return {"tool": name, "ok": False, "error": str(e)}
def tool_calculator(expression: str) -> dict:
allowed = set("0123456789+-*/(). ")
if not set(expression) <= allowed:
raise ValueError("表达式包含非法字符")
value = eval(expression, {"__builtins__": {}}, {})
return {"expression": expression, "value": value}
def tool_weather(city: str) -> dict:
return {"city": city, "temperature": 24, "condition": "多云"}
def tool_search(query: str) -> dict:
return {
"query": query,
"results": [
{"title": "示例结果 1", "url": "https://example.com/1"},
{"title": "示例结果 2", "url": "https://example.com/2"},
],
}
def build_default_registry() -> ToolRegistry:
reg = ToolRegistry()
reg.register(ToolSpec(
name="calculator",
description="计算数学表达式",
parameters={
"type": "object",
"properties": {"expression": {"type": "string"}},
"required": ["expression"],
},
handler=tool_calculator,
risk_level="low",
))
reg.register(ToolSpec(
name="weather",
description="查询城市天气",
parameters={
"type": "object",
"properties": {"city": {"type": "string"}},
"required": ["city"],
},
handler=tool_weather,
))
reg.register(ToolSpec(
name="search",
description="搜索信息",
parameters={
"type": "object",
"properties": {"query": {"type": "string"}},
"required": ["query"],
},
handler=tool_search,
))
return regapp/graph.py:
from typing import Callable, Dict, List, Optional, Any
from app.state import AgentState
NodeFn = Callable[[AgentState], AgentState]
RouterFn = Callable[[AgentState], str]
class Node:
def __init__(self, name: str, fn: NodeFn, kind: str = "function"):
self.name = name
self.fn = fn
self.kind = kind
def run(self, state: AgentState) -> AgentState:
return self.fn(state)
class Graph:
"""
极简图编排引擎。
- 支持普通边
- 支持条件边
- 支持并行节点
- 支持检查点(通过 checkpoint_store)
- 支持循环(通过条件边回到上游)
"""
def __init__(self, name: str = "graph"):
self.name = name
self.nodes: Dict[str, Node] = {}
self.edges: Dict[str, List[str]] = {}
self.conditional_edges: Dict[str, tuple] = {}
self.entry: Optional[str] = None
self.finish: str = "END"
self.checkpoint_store = None
self.parallel_targets: Dict[str, List[str]] = {}
def add_node(self, name: str, fn: NodeFn, kind: str = "function"):
self.nodes[name] = Node(name, fn, kind)
return self
def set_entry(self, name: str):
self.entry = name
return self
def add_edge(self, src: str, dst: str):
self.edges.setdefault(src, []).append(dst)
return self
def add_conditional_edge(
self,
src: str,
router: RouterFn,
mapping: Dict[str, str],
):
self.conditional_edges[src] = (router, mapping)
return self
def add_parallel(self, src: str, targets: List[str], join: str):
"""
并行:执行 targets 中所有节点,然后汇聚到 join。
"""
self.parallel_targets[src] = targets
for t in targets:
self.add_edge(t, join)
return self
def _save_checkpoint(self, state: AgentState, node: str):
if self.checkpoint_store:
self.checkpoint_store.save(state, node)
def run(
self,
state: AgentState,
resume_from: Optional[str] = None,
) -> AgentState:
current = resume_from or self.entry
if not current:
raise ValueError("图未设置入口节点")
visited = []
while current != self.finish:
node = self.nodes.get(current)
if not node:
state.errors.append(f"节点不存在:{current}")
state.status = "failed"
return state
visited.append(current)
try:
state = node.run(state)
except Exception as e:
state.errors.append(f"节点 {current} 异常:{e}")
state.status = "failed"
self._save_checkpoint(state, current)
return state
state.next_node = None
self._save_checkpoint(state, current)
if state.status in ("failed", "waiting_human"):
return state
# 并行分支
if current in self.parallel_targets:
targets = self.parallel_targets[current]
for t in targets:
t_node = self.nodes.get(t)
if t_node:
state = t_node.run(state)
self._save_checkpoint(state, t)
# 并行后走普通边
nexts = self.edges.get(current, [])
current = nexts[0] if nexts else self.finish
continue
# 条件边优先
if current in self.conditional_edges:
router, mapping = self.conditional_edges[current]
key = router(state)
nxt = mapping.get(key)
if not nxt:
state.errors.append(f"条件边未匹配:{current} -> {key}")
state.status = "failed"
return state
current = nxt
continue
# 普通边
nexts = self.edges.get(current, [])
if not nexts:
current = self.finish
else:
current = nexts[0]
if state.status == "running":
state.status = "done"
return state
def mermaid(self) -> str:
"""
输出 Mermaid 图,便于可视化。
"""
lines = ["graph TD"]
for src, dsts in self.edges.items():
for d in dsts:
lines.append(f" {src} --> {d}")
for src, (_, mapping) in self.conditional_edges.items():
for label, dst in mapping.items():
lines.append(f" {src} -- {label} --> {dst}")
return "\n".join(lines)app/nodes.py:
from app.state import AgentState
from app.loop import AgentLoop
def node_classify(state: AgentState) -> AgentState:
goal = state.goal
if any(k in goal for k in ["计算", "天气", "搜索"]):
state.set("task_type", "tool")
else:
state.set("task_type", "chat")
state.add_message("system", f"任务分类:{state.get('task_type')}")
return state
def node_plan(state: AgentState) -> AgentState:
state.set("plan", [
"识别目标",
"调用工具",
"汇总结果",
])
state.add_message("system", "已生成计划")
return state
def node_execute_loop(tools) -> callable:
def _run(state: AgentState) -> AgentState:
loop = AgentLoop(tools)
return loop.run(state)
return _run
def node_chat(state: AgentState) -> AgentState:
state.add_message("assistant", f"普通对话回复:{state.goal}")
state.status = "done"
return state
def node_review(state: AgentState) -> AgentState:
if state.errors:
state.set("review", "has_errors")
else:
state.set("review", "ok")
return state
def node_parallel_a(state: AgentState) -> AgentState:
state.set("branch_a", "A 完成")
return state
def node_parallel_b(state: AgentState) -> AgentState:
state.set("branch_b", "B 完成")
return state
def node_join(state: AgentState) -> AgentState:
state.set("join_result", f"{state.get('branch_a')} + {state.get('branch_b')}")
state.add_message("system", "并行分支已汇聚")
return state
def route_by_task(state: AgentState) -> str:
return state.get("task_type", "chat")
def route_by_review(state: AgentState) -> str:
return state.get("review", "ok")app/checkpoint.py:
import json
from datetime import datetime
from typing import Optional
from sqlmodel import SQLModel, Field, Session, select
from app.state import AgentState
from app.db import engine
class Checkpoint(SQLModel, table=True):
id: Optional[int] = Field(default=None, primary_key=True)
run_id: str = Field(index=True)
node: str
state_json: str
created_at: datetime = Field(default_factory=datetime.utcnow)
class SQLCheckpointStore:
def save(self, state: AgentState, node: str):
with Session(engine) as db:
cp = Checkpoint(
run_id=state.run_id,
node=node,
state_json=json.dumps(state.to_dict(), ensure_ascii=False, default=str),
)
db.add(cp)
db.commit()
def latest(self, run_id: str) -> Optional[dict]:
with Session(engine) as db:
row = db.exec(
select(Checkpoint)
.where(Checkpoint.run_id == run_id)
.order_by(Checkpoint.id.desc())
).first()
if not row:
return None
return {
"node": row.node,
"state": json.loads(row.state_json),
}
def history(self, run_id: str) -> list:
with Session(engine) as db:
rows = db.exec(
select(Checkpoint)
.where(Checkpoint.run_id == run_id)
.order_by(Checkpoint.id.asc())
).all()
return [
{"node": r.node, "state": json.loads(r.state_json), "at": str(r.created_at)}
for r in rows
]app/db.py:
from sqlmodel import SQLModel, create_engine, Session
DATABASE_URL = "sqlite:///./agent_graph.db"
engine = create_engine(
DATABASE_URL,
connect_args={"check_same_thread": False},
)
def init_db():
SQLModel.metadata.create_all(engine)
def get_session():
with Session(engine) as session:
yield session恢复执行:
def resume_run(run_id: str, graph, store: SQLCheckpointStore) -> AgentState:
latest = store.latest(run_id)
if not latest:
raise ValueError("没有找到检查点")
state = AgentState.from_dict(latest["state"])
state.status = "running"
return graph.run(state, resume_from=latest["node"])app/hitl.py:
from app.state import AgentState
def node_high_risk_action(state: AgentState) -> AgentState:
"""
高风险动作节点:需要人工确认。
"""
state.set("pending_action", {
"type": "delete_resource",
"target": "prod-db-01",
"reason": "用户请求清理",
})
state.status = "waiting_human"
state.add_message("system", "等待人工确认高风险操作")
return state
def approve(state: AgentState, approved: bool, comment: str = "") -> AgentState:
if state.status != "waiting_human":
return state
state.set("human_decision", {
"approved": approved,
"comment": comment,
})
if approved:
state.add_message("system", "人工已批准,继续执行")
state.status = "running"
else:
state.add_message("system", "人工已拒绝,终止流程")
state.status = "failed"
state.errors.append("人工拒绝")
return stateapp/subgraph.py:
from app.graph import Graph
from app.state import AgentState
def as_node(subgraph: Graph, name: str):
"""
把子图封装成节点,实现嵌套编排。
"""
def _run(state: AgentState) -> AgentState:
return subgraph.run(state)
return name, _run
def build_research_subgraph(tools) -> Graph:
from app.nodes import node_execute_loop
sub = Graph(name="research")
sub.add_node("search", node_execute_loop(tools))
sub.add_node("summarize", _summarize)
sub.set_entry("search")
sub.add_edge("search", "summarize")
sub.add_edge("summarize", "END")
return sub
def _summarize(state: AgentState) -> AgentState:
state.set("summary", "已汇总搜索结果")
return state多智能体示例:
def build_multi_agent_graph(tools) -> Graph:
from app.nodes import (
node_classify, node_plan, node_execute_loop,
node_chat, node_review, route_by_task, route_by_review,
)
g = Graph(name="multi_agent")
g.add_node("classify", node_classify)
g.add_node("plan", node_plan)
g.add_node("execute", node_execute_loop(tools))
g.add_node("chat", node_chat)
g.add_node("review", node_review)
g.set_entry("classify")
g.add_edge("classify", "plan")
g.add_conditional_edge("plan", route_by_task, {
"tool": "execute",
"chat": "chat",
})
g.add_edge("execute", "review")
g.add_conditional_edge("review", route_by_review, {
"ok": "END",
"has_errors": "plan", # 出错回到 plan,形成循环
})
return gapp/observability.py:
import json
from datetime import datetime
from typing import Any, Dict, List
from sqlmodel import SQLModel, Field, Session, select
from app.db import engine
class TraceEvent(SQLModel, table=True):
id: int | None = Field(default=None, primary_key=True)
run_id: str = Field(index=True)
node: str
event_type: str # node_start / node_end / tool_call / error
payload: str
created_at: datetime = Field(default_factory=datetime.utcnow)
class Tracer:
def __init__(self, run_id: str):
self.run_id = run_id
def log(self, node: str, event_type: str, payload: Dict[str, Any]):
with Session(engine) as db:
db.add(TraceEvent(
run_id=self.run_id,
node=node,
event_type=event_type,
payload=json.dumps(payload, ensure_ascii=False, default=str),
))
db.commit()
def timeline(self) -> List[dict]:
with Session(engine) as db:
rows = db.exec(
select(TraceEvent)
.where(TraceEvent.run_id == self.run_id)
.order_by(TraceEvent.id.asc())
).all()
return [
{
"node": r.node,
"type": r.event_type,
"payload": json.loads(r.payload),
"at": str(r.created_at),
}
for r in rows
]带追踪的节点包装:
def traced(name: str, fn, tracer: Tracer):
def _run(state: AgentState) -> AgentState:
tracer.log(name, "node_start", {"step": state.step})
try:
state = fn(state)
tracer.log(name, "node_end", {"status": state.status})
except Exception as e:
tracer.log(name, "error", {"error": str(e)})
raise
return state
return _runapp/main.py:
from fastapi import FastAPI, Depends, HTTPException
from pydantic import BaseModel
from sqlmodel import Session
from app.db import init_db, get_session
from app.state import AgentState
from app.tools import build_default_registry
from app.checkpoint import SQLCheckpointStore
from app.nodes import build_multi_agent_graph
from app.observability import Tracer
from app.hitl import approve
app = FastAPI(title="Agent Loop + Graph Engineer")
store = SQLCheckpointStore()
tools = build_default_registry()
@app.on_event("startup")
def on_startup():
init_db()
class RunRequest(BaseModel):
goal: str
user_id: int = 1
class ApproveRequest(BaseModel):
run_id: str
approved: bool
comment: str = ""
@app.get("/health")
def health():
return {"status": "ok"}
@app.post("/run")
def run_graph(payload: RunRequest):
graph = build_multi_agent_graph(tools)
graph.checkpoint_store = store
state = AgentState(user_id=payload.user_id, goal=payload.goal)
tracer = Tracer(state.run_id)
tracer.log("graph", "start", {"goal": payload.goal})
state = graph.run(state)
tracer.log("graph", "end", {"status": state.status})
return {
"run_id": state.run_id,
"status": state.status,
"messages": [m.__dict__ for m in state.messages],
"tool_results": state.tool_results,
"variables": state.variables,
"errors": state.errors,
"timeline": tracer.timeline(),
}
@app.get("/runs/{run_id}/checkpoints")
def get_checkpoints(run_id: str):
return store.history(run_id)
@app.post("/runs/resume")
def resume(run_id: str):
latest = store.latest(run_id)
if not latest:
raise HTTPException(status_code=404, detail="No checkpoint")
graph = build_multi_agent_graph(tools)
graph.checkpoint_store = store
state = AgentState.from_dict(latest["state"])
state.status = "running"
state = graph.run(state, resume_from=latest["node"])
return {"run_id": state.run_id, "status": state.status}
@app.post("/runs/approve")
def approve_run(payload: ApproveRequest):
latest = store.latest(payload.run_id)
if not latest:
raise HTTPException(status_code=404, detail="No checkpoint")
state = AgentState.from_dict(latest["state"])
state = approve(state, payload.approved, payload.comment)
if state.status == "running":
graph = build_multi_agent_graph(tools)
graph.checkpoint_store = store
state = graph.run(state, resume_from=latest["node"])
return {"run_id": state.run_id, "status": state.status}
@app.get("/graph/mermaid")
def graph_mermaid():
graph = build_multi_agent_graph(tools)
return {"mermaid": graph.mermaid()}tests/test_agent_graph.p
import os
import pytest
from fastapi.testclient import TestClient
os.environ["DATABASE_URL"] = "sqlite:///./test_agent_graph.db"
from app.main import app # noqa: E402
from app.graph import Graph # noqa: E402
from app.state import AgentState # noqa: E402
from app.tools import build_default_registry # noqa: E402
from app.loop import AgentLoop # noqa: E402
client = TestClient(app)
def test_health():
assert client.get("/health").json()["status"] == "ok"
def test_agent_loop_calculator():
tools = build_default_registry()
loop = AgentLoop(tools)
state = AgentState(goal="帮我计算 2+3*4")
state = loop.run(state)
assert state.status == "done"
assert any("value" in str(r) for r in state.tool_results)
def test_agent_loop_weather():
tools = build_default_registry()
loop = AgentLoop(tools)
state = AgentState(goal="查一下上海天气")
state = loop.run(state)
assert state.status == "done"
assert any("上海" in str(r) for r in state.tool_results)
def test_graph_run_tool_task():
r = client.post("/run", json={"goal": "帮我计算 2+3*4"})
assert r.status_code == 200
data = r.json()
assert data["status"] == "done"
assert len(data["tool_results"]) >= 1
def test_graph_run_chat_task():
r = client.post("/run", json={"goal": "你好,介绍一下你自己"})
assert r.status_code == 200
data = r.json()
assert data["status"] == "done"
def test_checkpoints_saved():
r = client.post("/run", json={"goal": "查一下上海天气"})
run_id = r.json()["run_id"]
r = client.get(f"/runs/{run_id}/checkpoints")
assert r.status_code == 200
assert len(r.json()) >= 2
def test_resume_run():
r = client.post("/run", json={"goal": "帮我计算 2+3*4"})
run_id = r.json()["run_id"]
r = client.post("/runs/resume", params={"run_id": run_id})
assert r.status_code == 200
assert r.json()["status"] in ("done", "failed")
def test_mermaid_output():
r = client.get("/graph/mermaid")
assert r.status_code == 200
assert "graph TD" in r.json()["mermaid"]
def test_conditional_routing():
from app.nodes import node_classify, route_by_task
state = AgentState(goal="帮我计算 1+1")
state = node_classify(state)
assert route_by_task(state) == "tool"
state2 = AgentState(goal="随便聊聊天")
state2 = node_classify(state2)
assert route_by_task(state2) == "chat"运行:
pytest -vAgent Loop 与 Graph Engineer 的关系:
Loop = 微观执行(一个节点内的推理循环)
Graph = 宏观编排(多个节点间的状态流转)核心公式:
Agent System = Graph(State, Nodes, Edges, Checkpoints, Loop)三条工程原则:
适用场景:
场景 | 推荐 |
|---|---|
简单工具调用 | 纯 Agent Loop |
多步业务流程 | Graph 编排 |
高风险操作 | Graph + HITL |
多智能体协作 | Graph + 子图 |
长任务 | Graph + 检查点 |
需要调试 | Graph + 追踪 |
当你能把业务逻辑画成一张图,并让每个节点内部跑一个可控的 Agent Loop 时,你就真正掌握了下一代智能体系统的工程方法。
agent-graph-engine/
├── app/
│ ├── __init__.py
│ ├── main.py
│ ├── db.py
│ ├── state.py
│ ├── loop.py
│ ├── tools.py
│ ├── graph.py
│ ├── nodes.py
│ ├── checkpoint.py
│ ├── hitl.py
│ ├── subgraph.py
│ ├── observability.py
│ └── llm.py
├── tests/
│ └── test_agent_graph.py
├── requirements.txt
└── Dockerfile运行顺序:
pip install -r requirements.txt
pytest -v
uvicorn app.main:app --reload原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。