LangGraph 研究智能体工程实践

简介: 先看一个真实问题:用户给出一个需要多来源调查的问题,系统如何拆分问题、查找资料、保存证据、处理冲突,并生成一份可以回溯的报告? 本文带大家build一个 Research Agent。它提供 Web 界面、HTTP API 和 CLI,…

💡提示:这篇是多平台同步发布,架构图、Mermaid 流程图、代码高亮、表格可能会出现排版错乱、丢失。

想要看完整图表、格式化代码、完整目录,请用电脑端访问原文阅读: https://blog.csdn.net/Maolei_NyaRu_/article/details/164214754

先看一个真实问题:用户给出一个需要多来源调查的问题,系统如何拆分问题、查找资料、保存证据、处理冲突,并生成一份可以回溯的报告?

本文带大家build一个 Research Agent。它提供 Web 界面、HTTP API 和 CLI,研究过程使用 LangGraph 编排,业务数据和检查点保存到 PostgreSQL。完整代码见 GitHub:

https://github.com/NyaRu-Kiss/SDD-Research_Agent

0. 先看成品

下面以这个研究任务为例:

比较 2024 年 Q3 中美 AI 芯片出口管制政策的影响。

图一:成品截图展示 项目采用前后端分离结构:backend/ 负责 API、工作流、数据库和外部服务,frontend/ 负责任务列表、详情、计划编辑和实时 trace,spec/ 保存产品与技术规格,sql/ 保存可审阅的数据库结构。运行方式和环境变量见 README.md。

1. Agent 是什么

1.1 一个能指导设计的定义

Agent 接收一个目标,在运行时根据当前状态选择下一步行动。这个行动可能是调用模型、调用搜索工具、读取网页、写入数据库、判断证据是否足够,或者向用户请求干预。

因此,使用了大模型的程序不一定是 Agent。一次调用模型并返回答案,更接近普通的问答函数。Agent 至少需要面对一个开放目标,并在多个步骤之间保存状态、使用工具、根据结果调整路径。

Research Agent 的输入是一个研究问题,中间状态包括研究计划、问题依赖、待处理问题、候选来源、证据、结论和预算。输出包括报告、来源、证据关联、进度和执行记录。

1.2 Agent 和固定工作流的区别

固定工作流适合步骤稳定、输入输出可预测的场景。Agent 适合目标明确但路径不确定的场景。

维度 固定工作流 Agent
下一步 按预先写好的顺序执行 根据当前状态和工具结果选择
分支 条件通常提前列出 可能因证据缺口、冲突或用户意见改变
循环 次数和边界较明确 围绕目标迭代,直到满足停止条件
外部世界 输入相对稳定 工具可能超时、返回空结果或互相矛盾
人的介入 通常在流程外处理 可以修改计划、暂停任务或要求继续调查
审计 记录步骤即可 还要记录决策、证据、重试和状态变化

生产系统通常是两者的组合:用工作流约束状态、权限和终止条件,用 Agent 决定研究什么、何时继续以及如何处理新信息。本项目正是这种组合。LangGraph 负责图的拓扑和状态边界,select_work、证据评估和收敛逻辑负责选择路径。

1.3 统一几个术语

  • ModelProvider:LLM 接口,负责模型调用和结构化输出。
  • SearchProvider:搜索接口,例如 Tavily。
  • ReaderProvider:网页阅读接口,例如 Firecrawl。
  • ExternalProviderToolProvider:对这些外部服务的统称。

1.4 什么时候不该用 Agent

如果步骤、输入和输出完全固定,普通函数或确定性工作流更容易测试,也更容易估算成本。只有一次问答、没有跨步骤状态时,也不必引入 LangGraph、检查点和一整套任务生命周期。

2. 需求和架构:先明确边界,再设计状态

2.1 这个项目需要解决什么问题

Research Agent 的需求可以压缩成五件事:

  1. 把复杂问题拆成可以独立调查、又能表达依赖关系的子问题。
  2. 让搜索、网页阅读、证据提取和评估组成研究循环,直到达到证据或资源边界。
  3. 把事实、分析和不确定结论区分开。
  4. 让重要结论可以追溯到来源和证据。
  5. 允许用户在计划阶段修改方向,在执行阶段中断并恢复任务。

2.2 当前 MVP 的边界

这点需要在架构图旁边直接说清楚:当前版本是单用户、单任务主线,没有登录、多用户协作、跨任务长期记忆和定时研究。并发只用于同一个任务中彼此独立的问题,数据库会话、模型调用和搜索调用仍受资源上限约束。

明确边界有两个好处。读者不会把示例误解成一个已经解决了多租户隔离的 SaaS;实现者也不会为了“架构完整”提前引入当前用不到的复杂度。

2.3 分层架构

图二:Research Agent 分层架构。箭头表示调用或依赖方向;Provider 只负责外部服务,Repository 负责持久化。

入口层 apicli 负责参数校验、响应模型和错误映射;application 负责一次完整的任务用例;domain 保存来源、证据和状态等不依赖具体框架的规则;workflows 将研究过程编排成图;infrastructure 隔离数据库、Repository 和外部 Provider。

这样的分层让“研究应该怎么走”和“搜索服务的 HTTP 请求怎么发”保持距离。测试工作流时可以注入 fake Provider,替换供应商时也不需要重写路由。

2.4 先定义 ResearchState

图三:ResearchState 状态结构。图状态保存路由所需的轻量字段,正文和历史记录保存在 PostgreSQL。

Agent 的工作记忆应该先于 Prompt 被设计出来。项目在 state.py 中定义 ResearchState,其中最重要的字段包括:

字段 用途
task_id / plan_id 定位任务和当前计划
pending_question_ids 等待研究的问题
dependencies 问题之间的前置关系
current_question_id 当前节点正在处理的问题
current_source_ids / current_evidence_ids 本轮来源和证据引用
interrupt_requested 是否在下一个安全边界暂停
route_reason 解释为什么选择某条路径

状态只保存路由所需的轻量数据和数据库记录 ID;网页正文、完整报告、历史事件则落 PostgreSQL。这样做可以避免把长文本复制到每个并行分支,也让恢复时有稳定的持久化来源。

源码中的状态定义还显式保存预算、循环计数和重试计数:

class ResearchState(TypedDict, total=False):
task_id: Annotated[str, latest]
phase: Annotated[Literal["规划", "研究", "汇总", "中断", "失败"], latest]
pending_question_ids: Annotated[list[str], latest]
dependencies: Annotated[dict[str, list[str]], merge_mapping]
budget: Annotated[dict[str, int], latest]
loop_counters: Annotated[dict[str, int], latest]
loop_limits: Annotated[dict[str, int], latest]
graph_steps: Annotated[int, latest]
retries_by_call: Annotated[dict[str, int], latest]

完整定义:state.py。

3. 用 LangGraph 搭出研究循环

3.1 LangGraph 的最小必要知识

读下面的代码只需要知道四个概念:StateGraph 用一个状态类型连接节点;conditional edge 根据状态返回下一条边;Send 可以把独立工作分发成并行分支;checkpointer 按 thread_id 保存状态,使任务能够恢复。

3.2 注册节点和连接主路径

图四:LangGraph 研究主流程。select_work 负责在执行、干预、中断和报告之间路由。

def build_research_graph(nodes=None, checkpointer=None):
graph = StateGraph(ResearchState)

for name in (
"load_task", "create_plan", "select_work",
"fan_out_questions", "search_sources", "read_sources",
"extract_evidence", "assess_evidence", "record_conclusion",
"fan_in_questions", "apply_intervention", "interrupt",
"generate_report", "fail_safely",
):
graph.add_node(name, nodes.get(name, passthrough))

graph.add_edge(START, "load_task")
graph.add_edge("create_plan", "select_work")
graph.add_edge("search_sources", "read_sources")
graph.add_edge("read_sources", "extract_evidence")
graph.add_edge("extract_evidence", "assess_evidence")
graph.add_edge("assess_evidence", "record_conclusion")
graph.add_edge("record_conclusion", "fan_in_questions")
graph.add_edge("fan_in_questions", "select_work")
graph.add_edge("generate_report", END)
return graph.compile(checkpointer=checkpointer)

完整实现:research_graph.py。

真实代码还为 load_taskselect_work 配置条件边,并保留可注入节点。节点注入很重要:测试可以只替换某个节点,不需要启动模型和数据库服务。

3.3 让状态决定下一步

_select_next 是一段很小但很关键的纯函数:

def _select_next(state):
if state.get("interrupt_requested"):
return "interrupt"
if state.get("pending_intervention_ids"):
return "intervention"
if state.get("current_question_id"):
return "search"

pending = state.get("pending_question_ids", [])
if pending:
return "fan_out" if len(state.get("current_question_ids", pending)) > 1 else "search"
return "report"

完整实现:research_graph.py。

这里的优先级是有意设计的:先响应暂停,再处理用户意见;如果已有当前问题就继续研究;没有当前问题时,从待处理问题中挑选单个或一批;所有问题完成后才生成报告。

把路由写成纯函数,可以单独测试“给定这个状态,下一步应该是什么”,而不需要真的调用搜索服务。

3.4 并行问题和依赖关系

图五:研究问题的并行与合流。无依赖问题可以同时执行,依赖前置结论的问题在合流后进入队列。

多个研究问题不必全部串行。例如“分别介绍三个框架”可以并行,“综合比较三个框架”则必须等待前面的资料完成。项目使用 LangGraph 的 Send

def _fan_out(state):
ready = state.get("current_question_ids", [])
return [
Send(
"search_sources",
{
   
**state,
"current_question_id": question_id,
"current_question_ids": [question_id],
"candidate_urls": [],
"current_source_ids": [],
},
)
for question_id in ready
]

完整实现:research_graph.py。

每个分支都有自己的当前问题和临时来源列表,完成后再回到 fan_in_questions。生产节点还用锁串行化共享的 SQLAlchemy AsyncSession,避免图上的并发分支同时操作同一个会话。

3.5 Prompt 工程:把自由文本变成受约束的决策

“先设计状态”不代表不需要 Prompt。Prompt 应该围绕状态和节点职责分层:系统 Prompt 规定证据边界、输出格式和禁止事项;节点 Prompt 只提供当前问题、任务约束和必要上下文。

先定义模型输出的边界。SearchQueryOutput 限制搜索词和域名数量;EvidenceOutput 要求模型返回可验证的原文引用;ReportClaimOutput 只允许报告 claim 指向已经保存的 conclusion_id

class SearchQueryOutput(BaseModel):
query: str = Field(min_length=1, max_length=500)
domains: list[str] = Field(default_factory=list, max_length=10)


class EvidenceOutput(BaseModel):
statement: str = Field(min_length=1, max_length=1200)
evidence_type: str = Field(pattern="^(事实|分析)$")
supporting_quote: str = Field(min_length=1, max_length=2000)
relevance_reason: str = Field(min_length=1, max_length=500)


class ReportClaimOutput(BaseModel):
conclusion_id: str = Field(min_length=1)
content: str = Field(min_length=1, max_length=2000)

完整定义:model_schemas.py。EvidenceOutput 的完整版本还包含引用片段位置;ReportOutput 还包含 answerkey_findingslimitationsclaims

ModelProvider.structured 将 schema 传给模型客户端。当前 OpenAI-compatible 配置使用 JSON mode,连接超时和输出解析错误会转成带 retryable=True 的领域错误:

def structured(self, schema: type[Any]) -> Any:
try:
if self._use_json_mode:
return self._client.with_structured_output(schema, method="json_mode")
return self._client.with_structured_output(schema)
except (ConnectionError, TimeoutError) as error:
raise ModelProviderError("模型服务暂不可用,可稍后重试。", retryable=True) from error
except OutputParserException as error:
raise ModelProviderError("模型输出无法按预期结构解析,可重试。", retryable=True) from error

完整实现:model_provider.py。

计划节点的 Prompt 带着用户问题、任务约束、精确 JSON schema、质量检查和示例发给模型:

prompt = {
   
"task": "Create a bounded research plan in Chinese.",
"user_question": question,
"constraints": constraints,
"requirements": (
"Return valid JSON only. Create 1-5 atomic questions. "
"Every question must define evidence_targets and minimum_evidence. "
"depends_on may only reference earlier questions."
),
"quality_checks": [
"Every material aspect is covered exactly once.",
"Dependencies are minimal and point only to earlier positions.",
],
}

完整实现:planning_nodes.py。该文件在解析失败后添加 correction 指令并重试,最终通过 PlanOutput.model_validate 再次校验。

证据提取节点要求逐条陈述并绑定引用;证据评估节点提供支持、冲突和置信度输入,最终分类交给领域服务;报告节点只读取已持久化的结论。结构化输出仍会失败,但失败会被转换为可记录、可重试或可降级的事件。

3.6 先测试图,再接真实模型

对应测试文件是 test_research_graph.py。测试使用替身节点验证路由、并行、中断和报告终止;Provider 层再单独测试模型格式错误、超时和空结果。这样可以把“图走错了”和“外部服务挂了”分开诊断。

def test_graph_compiles_and_uses_task_id_as_thread_id():
graph = build_research_graph()
assert graph is not None
assert graph_config("task-id") == {
   "configurable": {
   "thread_id": "task-id"}}
assert graph.invoke({
   "task_id": "task-id", "pending_question_ids": []})["task_id"] == "task-id"

完整测试文件:test_research_graph.py。

4. 研究计划:把模糊目标变成可执行问题

4.1 计划应该包含什么

计划至少要回答四个问题:研究目标是什么,要回答哪些子问题,子问题之间有什么依赖,最终报告需要覆盖哪些范围和输出形式。后续路由和进度统计都读取这份计划。

4.2 create_plan 的实现

planning_nodes.py 负责读取任务和创建计划,ModelPlanGenerator 调用模型,model_schemas.py 负责解析结构化结果。实现时先校验问题 ID 和依赖关系,再持久化计划;计划落库成功后才进入执行阶段。

伪代码如下:

async def create_plan(state, tasks, generator):
task = await tasks.get_task(state["task_id"])
plan = await generator.generate(task.question, task.constraints)
validate_question_ids(plan.questions)
validate_dependencies(plan.questions)
saved = await tasks.create_plan(task.id, plan)
return {
   "plan_id": str(saved.id), "pending_question_ids": question_ids(plan)}

完整实现:planning_nodes.py。

计划应保留版本和历史。用户修改方向后,系统需要定位当前执行使用的版本;任务恢复时,也需要校验计划与已完成问题是否匹配。

4.3 计划阶段的用户干预

这里的干预只指执行前的研究方向调整:修改问题、增加研究方向、改变时间范围、排除来源类型。计划更新通过 schema 校验,干预则作为待处理事件进入状态,由 apply_intervention 在安全边界应用。

执行中的暂停和继续放到第 8 章;HTTP API 如何暴露这两类操作放到第 9 章。这样“业务规则”和“接口形状”不会在三个章节重复出现。

5. 工具调用、预算和错误边界

5.1 Provider 只处理外部服务

model_provider.pytavily_provider.pyfirecrawl_provider.py 分别封装模型、搜索和网页阅读。节点决定什么时候调用、如何过滤结果、如何写入状态;Provider 不应该知道研究任务的完整业务流程。

搜索服务因此可以替换,测试也可以使用不联网的 fake。重试策略集中在应用侧,HTTP 客户端只处理自己的请求和响应。

5.2 搜索和网页阅读

search_sources 先生成受约束的查询,再应用任务中的域名、来源类型、排除列表和时间范围。它会去重候选 URL,并记录过滤原因。

read_sources 逐个读取候选页面,保存规范化 URL、标题、发布者、访问时间和内容哈希:

async def read_sources(state, tools, store):
source_ids = []
for url in state.get("candidate_urls", []):
page = await tools.read_webpage(url)
source = await store.save_source(
task_id=state["task_id"],
url=url,
canonical_url=canonicalize_url(str(page.canonical_url)),
title=str(getattr(page, "title", "")),
publisher=str(getattr(page, "publisher", "")),
accessed_at=page.accessed_at,
content=str(page.content),
content_hash=content_hash(str(page.content)),
)
source_ids.append(str(source.id))
return {
   "current_source_ids": source_ids}

保存内容哈希可以识别重复来源;保存访问时间则让报告引用具备基本的时间上下文。

完整实现:research_nodes.py。

5.3 预算是状态的一部分

研究 Agent 需要预算,运行边界才可预测。项目至少限制:

  • 搜索轮数和每轮来源数;
  • 模型 token、调用次数和重试次数;
  • 单任务运行时间、并发数和网页内容长度。

项目通过 settings 和 ResearchContext 传入 MAX_* 限制,再由 convergence_action 判断是否继续。达到资源上限后,系统选择 FINALIZE_WITH_GAPS,生成带有证据缺口的结果并停止搜索。

预算判断集中在一个函数中:

def convergence_action(state: ResearchState, context: ResearchContext) -> ConvergenceAction:
if state.get("graph_steps", 0) >= context.max_graph_steps:
return ConvergenceAction.FINALIZE_WITH_GAPS
if state.get("question_batches", 0) >= context.max_question_batches:
return ConvergenceAction.FINALIZE_WITH_GAPS
if any(rounds >= context.max_search_rounds_per_question
for rounds in state.get("search_rounds_by_question", {
   }).values()):
return ConvergenceAction.FINALIZE_WITH_GAPS
if any(retries >= context.max_retries_per_call
for retries in state.get("retries_by_call", {
   }).values()):
return ConvergenceAction.FINALIZE_WITH_GAPS
return ConvergenceAction.CONTINUE

完整实现:state.py。搜索轮数和重试次数也在同一个状态模型中递增,节点只根据上下文做路由判断。

预算需要进入状态和 trace。用户应能区分两种不确定性:资料本身不足,或系统已达到搜索、token 或时间上限。

5.4 节点失败和工作流失败

错误至少分三类:

  1. 节点级失败:模型输出格式错误、一次搜索超时、单个网页读取失败。可以在节点内重试、跳过或替换结果。
  2. 工作流级失败:任务不存在、状态损坏、依赖循环、连续多轮没有进展。需要停止任务并保留失败原因。
  3. 数据级不确定:来源互相冲突、证据不足。这类结果写入“不确定”或“证据缺口”结论,不进入系统异常流程。

生产工作流在 production_workflow.py 中给模型和 Provider 调用加了统一包装。可重试错误会指数退避;不可重试错误直接进入失败路径;每次重试都会写入诊断事件。

for attempt in range(max_retries + 1):
try:
return await call()
except Exception as error:
retryable = getattr(error, "retryable", False)
if not retryable or attempt >= max_retries:
raise
await diagnostics.event(
node=node,
event_type="retry_scheduled",
developer_metadata={
   "attempt_no": attempt + 1, "retryable": True},
)
await asyncio.sleep(0.25 * (2**attempt))

完整实现:production_workflow.py。

5.5 降级、死循环和强制收敛

图六:错误边界与收敛决策。实线中的重试、预算和 fail_safely 已在当前仓库实现;备用 Provider 切换属于扩展策略。

当前仓库已经实现了重试、预算上限、fail_safely 和带证据缺口的收敛。备用 Provider 切换、退回搜索摘要、按“无新增证据轮数”收敛,是下一步可以加入的策略,当前版本未实现。

现有的失败收口代码只改变任务状态并保留错误码:

async def fail_safely(state, task_store, error_code):
task = await task_store.get_task(state["task_id"])
if task is None:
raise LookupError("研究任务不存在。")
task.status = TaskStatus.FAILED
return {
   "route_reason": "研究无法继续,已有成果已保留。", "last_error_code": error_code}

完整实现:control_nodes.py。

扩展时可加入以下策略:

  • 搜索 Provider 失败时,切换备用 Provider,或者退回已有搜索摘要,并降低证据质量标记。
  • 网页读取失败时保留搜索结果,但不把摘要当成完整网页证据。
  • 连续若干轮没有新增来源、证据或已完成问题时,触发强制收敛。
  • 检测循环依赖、重复查询和反复读取同一 URL;达到上限后进入 fail_safely 或带缺口的报告路径。

所有重试、降级、路由和终止原因都应进入 trace。否则用户只会看到“报告不完整”,开发者也无法判断是证据不足还是系统故障。

6. 证据链:让报告可以被验证

6.1 Evidence 的数据模型

图七:来源、证据、结论关系。报告 claim 通过 conclusion_id 回溯到结论,再回到支持它的证据和来源。

项目在 evidence.pyresearch_output.pyresearch_repository.py 中保存证据和结论。一条证据至少要关联研究问题、来源、陈述内容、证据类型、置信度和支持关系。

来源是“我读过什么”,证据是“来源中的哪条信息支持什么问题”,结论是“系统根据多条证据做出的判断”。把三者分开,报告才有可能实现可追溯引用。

6.2 提取和评估

研究节点先从网页内容提取候选证据,再调用领域服务 assess_evidence。领域服务不负责联网,它只接收标准化的证据评估结果,并决定结论属于事实、分析还是不确定。

简化后的记录逻辑如下:

async def assess_and_record(state, store, context):
evidence_ids = state.get("current_evidence_ids", [])
assessments = [
EvidenceAssessment(item, EvidenceType.FACT, 0.7, 0.9)
for item in evidence_ids
]
decision = assess_evidence(assessments)

if should_finalize_with_gaps(state, context):
decision = uncertain("已达资源上限,证据不足。")

conclusion = await store.create_conclusion(
task_id=state["task_id"],
question_id=state["current_question_ids"][0],
content=decision.decision_basis,
classification=decision.classification,
confidence=decision.confidence,
decision_basis=decision.decision_basis,
)
return {
   "conclusion_id": str(conclusion.id)}

这段代码要关注最后一道约束:资源耗尽时,系统强制保留“不确定”语义,禁止把推断升级成事实。

完整实现:research_nodes.py,领域评估规则见 evidence_service.py。

6.3 冲突处理

不同来源出现冲突时,不能只选择模型最后一次生成的答案。系统至少要比较来源权威性、信息是否直接回答问题、发布时间、多个来源之间是否相互支持。

如果冲突无法消除,报告保留不同观点并解释原因。预算耗尽时,报告直接写出“当前证据不足”。第 5 章的预算与强制收敛规则会在这一层留下用户可见的结果。

6.4 用领域测试固定规则

test_evidence_service.py 覆盖低质量来源、证据不足、无法解决的冲突,以及明显更强的冲突证据。测试这些规则时不需要真实模型,输入可以是固定的 EvidenceAssessment,输出则应稳定地落在事实、分析或不确定分类上。

def test_insufficient_evidence_produces_uncertain_conclusion():
decision = assess_evidence([evidence(0.5, 0.9)])
assert decision.classification is ConclusionClassification.UNCERTAIN

完整测试文件:test_evidence_service.py。

7. 报告生成:从证据到最终产物

7.1 为什么不一次性生成全文

一次性 Prompt 让模型同时负责结构、事实、引用和语言,很难控制遗漏,也很难在中途恢复。更实际的做法是让报告生成消费已经确认的结论,把“研究”和“写作”拆开。

7.2 分阶段生成

图八:报告生成管线。每个阶段都有自己的输入和检查点,避免让一次模型调用同时承担结构、事实和引用。

一个可控的报告管线可以分成四步:

  1. 生成大纲:根据研究问题和已确认结论,固定章节顺序和比较维度。
  2. 生成段落:按章节生成文字,每个事实句携带内部 claim/evidence ID。
  3. 审校:检查引用覆盖、事实/推断标签、冲突和不确定性;失败时只重写具体章节。
  4. 渲染:根据 claim 与 evidence 关系注入来源标题、链接和访问时间,输出用户可读报告。

项目中的 ReportOutputReportClaimOutput 用 schema 约束报告结构,generate_report 负责在结论持久化后创建报告。引用由系统根据证据关系注入,不能让模型自由拼接 URL。

7.3 报告成功才关闭任务

control_nodes.pygenerate_report 先读取任务,再创建报告;报告成功落库后才把任务标记为 COMPLETED。如果报告生成失败,已有来源、证据和结论仍然保留,任务可以进入失败处理或稍后重试。

async def generate_report(state, task_store, store):
task = await task_store.get_task(state["task_id"])
if task is None:
raise LookupError("研究任务不存在。")
report = await store.create_report(
state["task_id"], {
   "conclusion_id": state.get("conclusion_id")}
)
task.status = TaskStatus.COMPLETED
return {
   "report_id": str(report.id)}

完整实现:control_nodes.py。

报告包含结论摘要、核心发现、比较或详细分析、限制与不确定性、来源清单。比较类问题先固定比较维度;决策类问题写出建议成立的前提和权衡;多个对象的介绍需要围绕这些维度组织。

7.4 UI 如何呈现可信度

TaskDetailPage.tsx 负责任务详情,ResearchTracePanel.tsx 负责过程事件和 trace。界面应该允许用户从结论跳到证据和来源,并明确标记事实、分析和不确定。

trace 对外展示节点名称、进度、输入输出摘要、重试和路由原因。它不包含模型思维链或不必要的内部推理文本。

8. 可恢复执行:中断、恢复和事件流

8.1 执行阶段的中断与恢复

control_nodes.py 中的 interrupt 职责很窄:读取任务、把状态设为中断,并返回路由原因。它不会删除已经收集的来源、证据和结论。

LangGraph 使用稳定的任务 ID 作为 thread_id

def graph_config(task_id: str):
return {
   "configurable": {
   "thread_id": task_id}}

完整实现:research_graph.py。

恢复时,ResearchService.resume 使用同一个配置重新运行图,checkpointer 读取之前的上下文。中断安排在节点边界,避免把数据库事务切在半途。

8.2 事件流和进度

图九:任务事件时序。终止事件同时发送给在线订阅者并写入缓存,重连客户端仍能获得最终状态。

task_events.py 定义任务事件发布器,API 通过 /{task_id}/events 提供 SSE。终止事件会缓存,因此客户端即使在任务完成后才重新连接,也能拿到最终状态。

发布器先缓存终止事件,再把事件复制到订阅者队列:

async def publish(self, event: TaskTerminalEvent) -> None:
async with self._lock:
self._latest[event.task_id] = event
subscribers = tuple(self._subscribers.get(event.task_id, set()))
for queue in subscribers:
if queue.full():
queue.get_nowait()
queue.put_nowait(event)

完整实现:task_events.py。

前端订阅事件时需要处理三件事:页面离开时取消订阅,收到终止事件后停止轮询,收到错误时回退到普通的进度查询。实时体验依赖这些状态边界和回退逻辑。

8.3 幂等、事务和审计

Repository 负责幂等写入和事务边界,PostgreSQL 同时保存业务数据和 LangGraph 检查点。节点生命周期统一经过 DiagnosticRecorder,记录开始、成功、失败、重试和路由原因。

恢复流程需要两项能力:系统能从检查点继续运行,并能解释当前状态的来源。

节点诊断包装器会在处理函数前后写入事件,异常时记录脱敏后的错误摘要:

async def node(self, node, state, handler):
await self.event(node=node, event_type="node_started", node_input=safe_data(state))
try:
result = await handler()
except Exception as error:
await self.event(
node=node,
event_type="node_failed",
error_summary=f"{type(error).__name__}: 节点执行失败。",
node_input=safe_data(state),
)
raise
return result

完整实现:diagnostics.py。完整版本还会记录完成事件、耗时、脱敏后的输出和错误类型。

9. 产品外壳:API、CLI 和前端

9.1 应用服务是唯一的用例入口

research_service.py 提供 createrunplanprogresssourcesevidencetracesreportsinterveneinterruptresume 等方法。

API 和 CLI 都调用应用服务:API 把 HTTP 请求转换成方法参数,CLI 把命令行参数转换成同样的方法调用。这样计划修改、中断恢复和报告查询不会出现两套逐渐分叉的业务逻辑。

计划修改对应 plan/intervene;执行暂停和继续对应 interrupt/resume。本节只讲操作如何暴露,具体规则分别见第 4 章和第 8 章。

9.2 HTTP API 资源

API 前缀是 /api/v1/research-tasks,资源可以按任务 ID 组织:

资源 典型操作
任务 创建、列表、详情
计划 查询、更新
执行 进度、SSE 事件、中断、恢复、干预
研究产物 来源、证据、trace、报告

异步启动接口返回任务标识,后续通过状态查询和事件流观察执行。请求和响应模型集中定义在 schemas.py,错误响应则由 research_routes.py 统一映射。

中断和恢复路由只做服务调用、提交事务和安排后台执行:

@router.post("/{task_id}/resume", status_code=202, response_model=TaskResponse)
async def resume(task_id, background_tasks, request, app_service=Depends(service)):
task = await app_service.resume(task_id)
await app_service.session.commit()
await request.app.state.task_events.clear(task.id)
background_tasks.add_task(schedule, app_service, task.id)
return as_response(task)

完整实现:research_routes.py。

9.3 前端研究控制台

前端由任务列表页、任务详情页、研究控制面板和 trace 面板组成。所有请求集中在 research.ts,页面组件不直接拼接后端 URL。

前端 API 客户端统一处理响应解析和错误码:

private async request<T>(path: string, init?: RequestInit): Promise<T> {
   
const response = await fetch(
`${
     this.baseUrl}/api/v1/research-tasks${
     path}`,
{
    headers: {
    "Content-Type": "application/json" }, ...init },
);
const text = await response.text();
let body: unknown = null;
try {
    body = text ? JSON.parse(text) : null; } catch {
    /* 保留通用错误 */ }
if (!response.ok) {
   
const error = body as {
    error?: {
    code?: string; message?: string } } | null;
throw new ApiError(error?.error?.code ?? "REQUEST_FAILED", error?.error?.message ?? "请求失败");
}
if (body === null) throw new ApiError("INVALID_RESPONSE", "服务返回了空响应");
return body as T;
}

完整实现:research.ts。

10. 测试策略:替身设计与运行门禁

10.1 测试伴随开发

工作流路由在第 3 章测试,证据规则在第 6 章测试,Provider、Repository、API 和前端测试应紧跟各自实现。最后集中测试,通常已经很难判断错误来自哪个边界。

10.2 外部依赖用替身隔离

ModelProviderSearchProviderReaderProvider 的 fake/double 固定以下场景:模型格式错误、Provider 超时、空搜索结果、网页读取失败、冲突证据和中断恢复。测试验证系统的决策逻辑;真实网络服务只用于少量集成验证。

模型替身只需要实现结构化调用接口:

class DeterministicModel:
def with_structured_output(self, schema):
return self

async def ainvoke(self, _input):
return {
   "query": "固定测试查询", "domains": []}

完整测试替身:test_model_provider.py。

10.3 运行门禁

后端运行:

cd backend
uv run ruff check .
uv run mypy app
uv run pytest -q

前端运行:

cd frontend
npm run lint
npm test
npm run build

需要 PostgreSQL 的测试使用独立测试库。开发服务测试完成后应停止对应终端中的进程,避免长期占用内存。

11. 复盘:几个容易踩的坑

  • 只写一个大 Prompt:缺少显式状态、结构化 schema 和停止条件。
  • 把搜索结果直接塞给模型:没有 URL 规范化、内容哈希和证据关系。
  • 所有问题串行执行:忽略 Send、问题依赖和并发边界。
  • 把推断写成事实:领域服务需要保留不确定性。
  • 没有预算和无进展检测:任务可能无限运行或静默失败。
  • 只展示最终报告:缺少进度、事件、trace 和恢复能力。
  • API、CLI、前端各实现一套业务逻辑:应用服务应成为唯一的用例入口。

这些问题分别对应状态、路由、预算、证据、恢复和应用服务边界。每个边界都需要明确的数据结构、失败路径和测试。

12. 结尾:从这个骨架继续扩展

这个项目的核心设计可以概括为:显式状态、受约束的工具、可验证的证据链、明确预算、可恢复执行和可审计事件。它们比某一个具体模型或 Prompt 更决定系统能否长期运行。

当前 MVP 仍然是单用户、单任务主线。后续可以沿着几个方向扩展:更细的来源质量模型、人工审核节点、缓存与成本控制、更多领域 Provider,以及面向多任务的权限和队列设计。

附录:运行与阅读索引

A. 关键文件

  • 工作流图:research_graph.py
  • 状态定义:state.py
  • 计划节点:planning_nodes.py
  • 生产节点:production_workflow.py
  • 研究节点:research_nodes.py
  • 证据规则:evidence_service.py
  • 应用服务:research_service.py
  • API 路由:research_routes.py
  • 前端详情页:TaskDetailPage.tsx
相关文章
人工智能 缓存 前端开发
12720 75
|
5天前
|
人工智能 自然语言处理 安全
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
本文聚焦阿里云2026年推出的三款自研AI办公产品,清晰拆解千问办公、Qoder Teams、Qoder CN的差异化定位与能力边界:千问办公主打职场全场景提效,支持自然语言指令一键完成PPT生成、数据分析等高频办公任务;Qoder Teams面向程序员团队,深度整合AI代码生成、团队协同与企业知识库能力;Qoder CN则专为金融、政务等强合规场景打造,实现数据不出境与VPC私有化部署。文章同步给出分场景选型指南与最新活动定价,帮助不同类型的企业按需组合产品,实现业务岗、研发岗与强合规场景的AI能力全覆盖。
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
Web App开发 人工智能 API
1605 2
|
人工智能 JavaScript 开发工具
DeepSeek Harness 本地安装与使用指南
DeepSeek Harness(DSH)是DeepSeek AI开源的Agent运行框架,支持本地文件操作、命令执行与工具调用。基于Cordis插件架构,具备高扩展性与强可控性,适合开发者搭建可控Agent环境或开展模型基准测试。当前为开发者预览版,需Node.js环境,推荐先用`npx @deepseek-ai/dsh web`快速体验。
4963 0
人工智能 Java BI
1709 1
人工智能 JavaScript 测试技术
2671 2
开发工具 Swift git
2014 6
人工智能 JavaScript 测试技术
1272 5

热门文章

最新文章