如果你正在探索如何构建一个真正可用的 AI Agent,而不仅仅是调用 API 返回一段文本,那么你很可能已经遇到了几个核心难题:如何让 Agent 记住上下文?如何管理多轮对话的复杂状态?如何在 Agent 执行出错或中断后,优雅地恢复现场,而不是让用户从头再来?
这正是 DeepSeek Harness 试图系统化解决的关键问题。它不是一个简单的聊天界面包装,而是一个面向生产环境的Agent 编排与状态管理框架。很多人初次接触时,可能会被 “Harness”(马具)这个名字误导,以为它只是个“套”在模型外的工具。但实际上,它的核心价值在于内部那套严谨的Turn(回合)、Step(步骤)和 Session(会话)状态机模型。这套模型决定了 Agent 如何思考、执行、回溯和重建,是理解其设计哲学和工程实践的关键。
本文将深入 DeepSeek Harness 源码,解析其 Agent 的核心运行机制。我们不会停留在表面 API 调用,而是聚焦于三个最易混淆但至关重要的概念:Turn、Step 与 Session 重建。通过源码分析,你将理解:
- Agent 的一次“思考-执行”循环是如何被拆解和管理的。
- 当 Agent 执行失败或用户想要回退时,框架如何实现精准的“状态回滚”或“会话恢复”。
- 如何借鉴其设计思想,应用到自己的 Agent 项目中,构建更健壮、可维护的智能体系统。
读完本文,你将能清晰地画出 Harness 中 Agent 的状态流转图,并能在自己的项目中实现类似的会话持久化与故障恢复机制。
1. 核心问题:为什么 Agent 需要复杂的“状态管理”?
在开始源码解析前,我们必须先达成一个共识:一个功能完备的 Agent 与一个简单的“问答函数”有本质区别。
想象一个场景:你让 Agent “帮我分析一下这个季度的销售数据,并生成一份报告”。一个简单的实现可能是:将整个提示词和文件扔给大模型,等待它一次性输出所有内容。这种方式存在明显问题:
- 不可控:如果生成报告到一半中断,所有中间状态丢失,必须重来。
- 难交互:你无法在它“分析数据”这一步完成后,先确认分析结果,再让它继续“撰写报告”。
- 难调试:当最终结果不符合预期时,你无法定位是“数据理解”、“分析逻辑”还是“报告格式”哪个环节出了问题。
DeepSeek Harness 的解决方案是将这个宏观任务分解为可管理的原子单元,并对每个单元的状态进行持久化跟踪。这就是Session、Turn、Step三层抽象的来源。
- Session(会话):对应一个最高层级的任务或一次用户对话的生命周期。它包含了完成这个任务所需的所有上下文、历史记录以及最终状态(成功、失败、进行中)。Session 重建就是指在服务重启或连接中断后,能够从持久化存储中恢复出这个完整的任务上下文,让 Agent 仿佛从未中断一样继续工作。
- Turn(回合):在 Session 内,一次明确的“用户输入 -> Agent 响应”的交互周期。例如,用户说“分析销售数据”,Agent 开始执行,这构成一个 Turn。一个复杂的任务可能包含多个 Turn(如:用户追问、Agent 请求澄清)。
- Step(步骤):这是最精细的粒度。在一个 Turn 内部,Agent 的“思考-执行”过程会被分解为多个 Step。例如,一个典型的 ReAct(Reasoning and Acting)模式可能包含:
Thought Step(思考下一步做什么)、Action Step(调用一个工具,如查询数据库)、Observation Step(获取工具执行结果)。每个 Step 都有其类型、输入、输出和状态。
理解了这三层模型,我们就能明白,Harness 的强大之处在于它对 Step 级状态的精细化管理。这使得“回到上一步”、“从错误步骤重试”、“并行执行多个步骤”等高级功能成为可能。接下来,我们从源码层面看这套机制如何运转。
2. 源码核心:状态定义与流转
我们首先关注定义这些状态的源码。在 Harness 的代码库中(通常位于src/core/或类似目录下),可以找到核心的状态类。
2.1 Session、Turn、Step 的实体定义
以下是根据 Harness 设计模式推断的简化版核心类定义,它清晰地展示了层级关系和关键属性:
# 文件路径:src/core/session.py (示例结构) from enum import Enum from typing import List, Optional, Dict, Any from datetime import datetime from pydantic import BaseModel class SessionStatus(Enum): CREATED = "created" RUNNING = "running" PAUSED = "paused" COMPLETED = "completed" FAILED = "failed" class Session(BaseModel): """会话实体,代表一个完整的任务生命周期。""" id: str # 会话唯一标识 user_id: Optional[str] # 关联用户 status: SessionStatus # 当前状态 context: Dict[str, Any] # 会话级上下文,如任务目标、全局变量 turns: List['Turn'] = [] # 包含的所有回合 created_at: datetime updated_at: datetime metadata: Dict[str, Any] = {} # 扩展元数据 def rebuild_context(self): """会话重建的核心方法:从持久化存储加载 turns 和 steps,恢复上下文。""" # 通常会从数据库加载关联的 turns 和 steps # 并可能根据最新的 step 状态,计算出当前的推理上下文 pass# 文件路径:src/core/turn.py class TurnStatus(Enum): INITIATED = "initiated" PROCESSING = "processing" WAITING_FOR_INPUT = "waiting_for_input" COMPLETED = "completed" ERROR = "error" class Turn(BaseModel): """回合实体,代表一次交互循环。""" id: str session_id: str # 所属会话 sequence: int # 在当前会话中的顺序 user_input: str # 用户输入 status: TurnStatus steps: List['Step'] = [] # 包含的所有步骤 started_at: datetime finished_at: Optional[datetime] # 可能包含该回合的最终输出或摘要 final_output: Optional[Dict[str, Any]]# 文件路径:src/core/step.py class StepType(Enum): THOUGHT = "thought" ACTION = "action" OBSERVATION = "observation" FINAL = "final" class StepStatus(Enum): PENDING = "pending" EXECUTING = "executing" SUCCESS = "success" FAILED = "failed" CANCELLED = "cancelled" class Step(BaseModel): """步骤实体,代表 Agent 执行的最小原子操作。""" id: str turn_id: str # 所属回合 step_type: StepType status: StepStatus content: Dict[str, Any] # 步骤内容,如思考内容、动作参数、观察结果 output: Optional[Dict[str, Any]] # 步骤执行输出 error_info: Optional[str] # 如果失败,错误信息 created_at: datetime # 用于步骤间关联,如前驱步骤ID parent_step_id: Optional[str] # 执行顺序或优先级 order: int关键点解析:
- 清晰的层级:
Session包含Turn,Turn包含Step。通过session_id和turn_id外键关联,构成一棵状态树。 - 状态枚举:每个实体都有明确的状态枚举。这是状态机流转的基础,也是实现“暂停”、“重试”、“恢复”等操作的前提。
- 内容与输出分离:
Step的content和output分离。content是输入(如“调用天气API,城市=北京”),output是执行结果(如“北京晴,25度”)。这便于记录和回放。 - 可扩展性:通过
context和metadata字段,可以附加任意自定义数据,适应不同业务场景。
2.2 状态机的驱动引擎:Orchestrator
定义了状态,还需要一个驱动状态流转的引擎。在 Harness 中,这个角色通常是Orchestrator(编排器)或AgentEngine。
# 文件路径:src/core/orchestrator.py (核心逻辑示例) class Orchestrator: def __init__(self, session_repository, llm_client, tool_registry): self.session_repo = session_repository self.llm = llm_client self.tools = tool_registry async def process_turn(self, session_id: str, user_input: str) -> Turn: """处理一个新的用户输入回合。""" # 1. 获取或创建Session session = await self.session_repo.get(session_id) if not session: session = Session(id=session_id, status=SessionStatus.CREATED, ...) await self.session_repo.save(session) # 2. 创建新的Turn new_turn = Turn( session_id=session.id, sequence=len(session.turns) + 1, user_input=user_input, status=TurnStatus.INITIATED, ... ) session.turns.append(new_turn) session.status = SessionStatus.RUNNING # 3. 核心:执行Step循环 (ReAct模式示例) await self._execute_step_loop(new_turn, session.context) # 4. 更新状态并保存 new_turn.status = TurnStatus.COMPLETED session.updated_at = datetime.now() await self.session_repo.save(session) return new_turn async def _execute_step_loop(self, turn: Turn, context: Dict): """执行一个Turn内的Step循环,直到产生最终答案或失败。""" max_steps = 10 for step_count in range(max_steps): # 决定下一步做什么 (由LLM根据上下文决定) next_step_plan = await self._plan_next_step(turn, context) step = Step(turn_id=turn.id, step_type=next_step_plan.type, ...) # 执行该步骤 step_output, step_status = await self._execute_single_step(step, next_step_plan, context) # 保存步骤结果 step.output = step_output step.status = step_status turn.steps.append(step) # 更新上下文(将本次步骤的观察结果加入) context.update(self._extract_observation_from_step(step)) # 检查是否应该结束(LLM决定或遇到最终答案) if await self._should_finish_turn(step, context): turn.final_output = self._compile_final_answer(turn.steps) break # 循环结束,可能成功也可能达到最大步数限制 if step_count == max_steps - 1: turn.status = TurnStatus.ERROR # 记录错误:可能陷入循环 async def _execute_single_step(self, step: Step, plan, context) -> (Dict, StepStatus): """执行单个步骤。""" try: if step.step_type == StepType.ACTION: tool_name = plan.content.get("tool") tool = self.tools.get(tool_name) if not tool: raise ValueError(f"Tool {tool_name} not found") # 调用工具 result = await tool.execute(**plan.content.get("parameters", {})) return {"observation": result}, StepStatus.SUCCESS elif step.step_type == StepType.THOUGHT: # 思考步骤可能只是记录LLM的推理,不产生外部作用 return {"thought": plan.content}, StepStatus.SUCCESS # ... 处理其他StepType except Exception as e: # 记录详细的错误信息 step.error_info = str(e) return {"error": str(e)}, StepStatus.FAILED引擎工作流程解读:
process_turn是入口,它管理 Turn 的生命周期。_execute_step_loop是核心,它实现了 ReAct 等循环模式。每一次循环都对应一个 Step 的创建、执行和持久化。_execute_single_step是原子操作执行器。这里会根据 Step 类型(如ACTION)调用相应的工具,并捕获异常。- 关键设计:每一步的结果(
step.output)都会更新到context中,作为后续步骤的输入。这就是 Agent “记忆”的来源。
3. 灵魂功能:Session 重建机制详解
“重建”是 Harness 应对故障和提供交互式体验的核心。其本质是从持久化存储(如数据库)中,按 Session ID 加载出完整的、包含所有 Turn 和 Step 的状态树,并重新初始化 Orchestrator 的上下文,使其能从中断点继续执行。
3.1 重建的触发场景
- 服务重启:服务器崩溃或更新后重启。
- 连接中断:用户客户端(如网页、APP)断开重连。
- 主动恢复:用户打开一个历史未完成的任务。
- 调试与回滚:开发者希望从某个特定步骤重新运行。
3.2 重建的核心逻辑
重建的核心在Session类的rebuild_context方法,以及 Orchestrator 如何利用重建后的 Session。
# 文件路径:src/core/session_repository.py (数据访问层示例) class SessionRepository: async def rebuild_session(self, session_id: str) -> Session: """从数据库重建完整的Session对象,包括其所有Turns和Steps。""" # 1. 加载Session基础信息 session_data = await self.db.fetch_one("SELECT * FROM sessions WHERE id = $1", session_id) if not session_data: raise NotFoundError(f"Session {session_id} not found") session = Session(**session_data) # 2. 加载该Session下的所有Turns(按sequence排序) turns_data = await self.db.fetch_all( "SELECT * FROM turns WHERE session_id = $1 ORDER BY sequence", session_id ) for turn_data in turns_data: turn = Turn(**turn_data) # 3. 加载每个Turn下的所有Steps(按order排序) steps_data = await self.db.fetch_all( "SELECT * FROM steps WHERE turn_id = $1 ORDER BY \"order\"", turn.id ) turn.steps = [Step(**data) for data in steps_data] session.turns.append(turn) # 4. 重新计算当前上下文 # 上下文通常是所有Step的output的某种聚合,或者是最后一个有效Step之后的“现场” session.context = self._reconstruct_context_from_steps(session.turns) # 5. 根据最后一个Step的状态,确定Session的当前状态 session.status = self._determine_session_status(session.turns) return session def _reconstruct_context_from_steps(self, turns: List[Turn]) -> Dict: """从步骤历史中重建出Agent继续执行所需的上下文。""" context = {} # 简化:将每个Step的output合并到context中,最新的覆盖旧的 for turn in turns: for step in turn.steps: if step.output and step.status == StepStatus.SUCCESS: context.update(step.output) # 注意:实际合并逻辑可能更复杂 return context def _determine_session_status(self, turns: List[Turn]) -> SessionStatus: """根据最后一个Turn和Step的状态推断Session状态。""" if not turns: return SessionStatus.CREATED last_turn = turns[-1] if last_turn.status == TurnStatus.ERROR: return SessionStatus.FAILED if last_turn.status == TurnStatus.COMPLETED: # 检查是否所有计划内的Turn都完成了 return SessionStatus.COMPLETED # 如果最后一个Turn还在处理中,或者处于等待输入状态,则Session是运行中或暂停 return SessionStatus.RUNNING # 或 PAUSED重建后的使用:当一个中断的连接重新建立,客户端会发送session_id。服务端会调用rebuild_session,然后将恢复的session对象交给一个Orchestrator实例。
# 文件路径:src/api/agent_controller.py (API层示例) @app.post("/agent/continue") async def continue_session(session_id: str, new_input: Optional[str] = None): # 1. 重建会话 session = await session_repo.rebuild_session(session_id) orchestrator = get_orchestrator() # 获取Orchestrator实例 # 2. 如果会话之前是等待输入状态,且提供了新输入,则继续处理 if session.status == SessionStatus.PAUSED and new_input: # 找到最后一个未完成的Turn,或者创建新的Turn last_turn = session.turns[-1] if last_turn.status == TurnStatus.WAITING_FOR_INPUT: last_turn.user_input = new_input last_turn.status = TurnStatus.PROCESSING # 使用重建的上下文继续执行_step_loop await orchestrator._execute_step_loop(last_turn, session.context) # 3. 如果只是查询状态,则直接返回重建后的会话信息 else: return {"session": session, "message": "Session restored."}3.3 重建的技术挑战与解决方案
- 上下文序列化:
context字典可能包含复杂对象(如数据库连接、API客户端)。重建时,通常只恢复纯数据部分,轻量级对象可以重新创建,重量级对象需通过工厂或依赖注入重新获取。 - 工具状态恢复:如果 Step 涉及调用一个有状态的外部工具(如启动了一个虚拟机),重建 Session 可能无法自动恢复该工具的状态。Harness 通常建议工具设计为幂等或提供状态查询接口。
- 并发安全:同一个 Session 被多个请求同时尝试恢复和继续执行时,需要加锁(如基于
session_id的分布式锁)防止状态混乱。 - 存储选择:需要支持快速查询嵌套结构。关系数据库(如 PostgreSQL)需要合理的表结构设计(
sessions,turns,steps三张表)。文档数据库(如 MongoDB)可能更自然,可以直接存储整个 Session 文档。
4. 实战:基于 Harness 设计思想实现一个简易 Agent 系统
理解了原理,我们可以抛开 Harness 框架本身,实现一个具备核心状态管理能力的小型 Agent 系统。我们将使用 Python 和 FastAPI,并用内存字典模拟持久化。
4.1 环境准备
# 创建项目目录并初始化 mkdir mini_agent_system && cd mini_agent_system python -m venv venv source venv/bin/activate # Windows: venv\Scripts\activate pip install fastapi uvicorn pydantic openai4.2 核心模型定义 (models.py)
# 文件路径:models.py from enum import Enum from typing import List, Optional, Dict, Any from pydantic import BaseModel from datetime import datetime class StepType(str, Enum): THOUGHT = "thought" ACTION = "action" OBSERVATION = "observation" FINAL = "final" class StepStatus(str, Enum): PENDING = "pending" EXECUTING = "executing" SUCCESS = "success" FAILED = "failed" class Step(BaseModel): id: str turn_id: str step_type: StepType status: StepStatus content: Dict[str, Any] output: Optional[Dict[str, Any]] = None error: Optional[str] = None created_at: datetime = datetime.now() class TurnStatus(str, Enum): INITIATED = "initiated" PROCESSING = "processing" COMPLETED = "completed" ERROR = "error" class Turn(BaseModel): id: str session_id: str sequence: int user_input: str status: TurnStatus steps: List[Step] = [] started_at: datetime = datetime.now() finished_at: Optional[datetime] = None final_output: Optional[str] = None class SessionStatus(str, Enum): CREATED = "created" RUNNING = "running" PAUSED = "paused" COMPLETED = "completed" FAILED = "failed" class Session(BaseModel): id: str status: SessionStatus context: Dict[str, Any] = {} # 存储对话历史、变量等 turns: List[Turn] = [] created_at: datetime = datetime.now() updated_at: datetime = datetime.now()4.3 状态存储与重建服务 (storage.py)
# 文件路径:storage.py from typing import Dict, Optional from models import Session, Turn, Step import uuid class InMemoryStorage: """简易内存存储,模拟数据库。生产环境需替换为真实数据库。""" def __init__(self): self.sessions: Dict[str, Session] = {} self.turns: Dict[str, Turn] = {} self.steps: Dict[str, Step] = {} # Session 相关操作 def create_session(self) -> Session: session_id = str(uuid.uuid4()) session = Session(id=session_id, status=SessionStatus.CREATED) self.sessions[session_id] = session return session def get_session(self, session_id: str) -> Optional[Session]: """获取Session,但不包含关联的Turns和Steps。""" return self.sessions.get(session_id) def rebuild_session(self, session_id: str) -> Optional[Session]: """重建Session:加载其所有Turns和Steps。""" session = self.get_session(session_id) if not session: return None # 清空现有 turns,准备重建 session.turns.clear() # 找出所有属于此 session 的 turns (模拟关联查询) session_turns = [t for t in self.turns.values() if t.session_id == session_id] # 按 sequence 排序 session_turns.sort(key=lambda x: x.sequence) for turn in session_turns: # 重建每个 turn 的 steps turn_steps = [s for s in self.steps.values() if s.turn_id == turn.id] turn_steps.sort(key=lambda x: x.created_at) turn.steps = turn_steps session.turns.append(turn) # 重建上下文:取最后一个成功的 Step 的 output all_steps = [step for turn in session.turns for step in turn.steps] successful_steps = [s for s in all_steps if s.status == StepStatus.SUCCESS and s.output] if successful_steps: # 简单策略:最后一个成功步骤的 output 作为上下文 session.context = successful_steps[-1].output or {} # 更新 session 状态 if session.turns: last_turn = session.turns[-1] if last_turn.status == TurnStatus.COMPLETED: session.status = SessionStatus.COMPLETED elif last_turn.status == TurnStatus.ERROR: session.status = SessionStatus.FAILED else: session.status = SessionStatus.RUNNING return session def save_session(self, session: Session): self.sessions[session.id] = session # 同时保存其下的 turns 和 steps (实际数据库会有事务) for turn in session.turns: self.turns[turn.id] = turn for step in turn.steps: self.steps[step.id] = step # 全局存储实例 storage = InMemoryStorage()4.4 核心编排引擎 (orchestrator.py)
# 文件路径:orchestrator.py import uuid from typing import Optional from models import Session, Turn, Step, StepType, StepStatus, TurnStatus, SessionStatus from storage import storage class MiniOrchestrator: def __init__(self, llm_client=None): # 简化,暂不集成真实LLM self.llm = llm_client async def process_user_input(self, session_id: str, user_input: str) -> Turn: """处理用户输入,创建或继续一个Turn。""" # 1. 尝试重建现有Session session = storage.rebuild_session(session_id) if not session: # 如果Session不存在,则创建新的 session = storage.create_session() session_id = session.id # 2. 创建新的Turn new_turn_id = str(uuid.uuid4()) new_turn = Turn( id=new_turn_id, session_id=session_id, sequence=len(session.turns) + 1, user_input=user_input, status=TurnStatus.PROCESSING ) session.turns.append(new_turn) session.status = SessionStatus.RUNNING session.updated_at = datetime.now() # 3. 模拟执行一个简单的两步流程:思考 -> 行动 # Step 1: 思考 thought_step = Step( id=str(uuid.uuid4()), turn_id=new_turn_id, step_type=StepType.THOUGHT, status=StepStatus.SUCCESS, content={"reasoning": f"用户说: {user_input}. 我需要理解他的意图。"}, output={"conclusion": "这是一个问候,需要友好回应。"} ) new_turn.steps.append(thought_step) # Step 2: 行动 (生成回复) action_step = Step( id=str(uuid.uuid4()), turn_id=new_turn_id, step_type=StepType.ACTION, status=StepStatus.SUCCESS, content={"action": "generate_response", "params": {"tone": "friendly"}}, output={"response": f"你好!我收到了你的消息:'{user_input}'。我是一个演示Agent,目前状态管理正常。"} ) new_turn.steps.append(action_step) # 4. 完成Turn new_turn.status = TurnStatus.COMPLETED new_turn.finished_at = datetime.now() new_turn.final_output = action_step.output.get("response") # 5. 更新Session上下文 (将最后输出加入上下文) session.context["last_response"] = new_turn.final_output session.context["last_user_input"] = user_input # 6. 保存所有状态 storage.save_session(session) return new_turn4.5 Web API 入口 (main.py)
# 文件路径:main.py from fastapi import FastAPI, HTTPException from pydantic import BaseModel from typing import Optional from orchestrator import MiniOrchestrator from storage import storage import uuid app = FastAPI(title="Mini Agent System") orchestrator = MiniOrchestrator() class UserRequest(BaseModel): session_id: Optional[str] = None # 不提供则创建新会话 message: str class SessionResponse(BaseModel): session_id: str status: str response: str turns_count: int @app.post("/chat", response_model=SessionResponse) async def chat_endpoint(request: UserRequest): # 确定session_id session_id = request.session_id or str(uuid.uuid4()) # 处理用户输入 turn = await orchestrator.process_user_input(session_id, request.message) # 重建session以获取最新状态 session = storage.rebuild_session(session_id) if not session: raise HTTPException(status_code=500, detail="Session not found after processing") return SessionResponse( session_id=session_id, status=session.status.value, response=turn.final_output or "No response generated.", turns_count=len(session.turns) ) @app.get("/session/{session_id}") async def get_session(session_id: str): """获取并重建一个会话的完整状态,用于前端恢复或调试。""" session = storage.rebuild_session(session_id) if not session: raise HTTPException(status_code=404, detail="Session not found") return session if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)4.6 运行与测试
- 启动服务:
uvicorn main:app --reload - 发起第一次请求(创建新会话):
响应会包含一个curl -X POST "http://127.0.0.1:8000/chat" \ -H "Content-Type: application/json" \ -d '{"message": "你好,Agent"}'session_id。 - 使用同一个 session_id 继续对话:
观察响应,Agent 的回复会体现上下文(虽然我们示例逻辑简单)。curl -X POST "http://127.0.0.1:8000/chat" \ -H "Content-Type: application/json" \ -d '{"session_id": "YOUR_SESSION_ID", "message": "还记得我吗?"}' - 查询会话完整状态:
你将看到完整的 Session 对象,包含所有 Turns 和 Steps,这就是Session 重建后返回的数据。curl "http://127.0.0.1:8000/session/YOUR_SESSION_ID"
5. 常见问题与排查思路
在实现和运行此类 Agent 系统时,你会遇到一些典型问题。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| Session 重建后上下文丢失 | 1.rebuild_session方法没有正确加载 Steps 的output。2. Steps 的 output字段在保存时为None或空。3. 上下文合并逻辑( _reconstruct_context_from_steps)有误。 | 1. 打印重建后 Session 的所有 Steps,检查output字段。2. 检查 _execute_single_step方法是否成功将结果赋值给step.output。3. 单步调试上下文合并函数。 | 1. 确保每个成功 Step 都将其结果写入output。2. 优化上下文合并策略,例如使用更复杂的结构(如列表)记录历史,而非简单覆盖。 3. 考虑将关键的上下文变量显式存储在 Session.context中。 |
| Agent 陷入循环,无限创建 Steps | 1. 结束条件判断(_should_finish_turn)逻辑不健全。2. LLM 的提示词没有引导其输出结束信号(如 Final Answer:)。3. 最大步数( max_steps)设置过高或未生效。 | 1. 检查循环日志,看 Step 的内容是否重复。 2. 检查 LLM 返回的文本是否包含预期的结束标记。 3. 确认 max_steps检查在循环中生效。 | 1. 强化结束判断逻辑,除了依赖 LLM,还可以加入超时、重复动作检测。 2. 在提示词中明确要求 LLM 在完成任务后输出特定结束语。 3. 设置合理的 max_steps(如 15-20),并在达到时强制结束,记录为“超时”错误。 |
| 工具(Action Step)执行失败后,Agent 不会处理错误 | 1._execute_single_step中捕获了异常但未更新 Step 状态为FAILED。2. Orchestrator 没有根据失败状态调整后续流程(如重试或终止)。 | 1. 检查 Step 执行后的status和error_info字段。2. 在 _execute_step_loop中检查上一步状态,并添加错误处理分支。 | 1. 确保异常被捕获,且step.status = StepStatus.FAILED和step.error_info被正确设置。2. 在循环中,如果上一个 Step 失败,可以决定是重试、询问用户还是终止 Turn。 |
| 并发请求导致同一个 Session 状态混乱 | 多个请求同时处理同一个session_id,同时读写其 Turns 和 Steps。 | 观察数据库或日志中,Steps 的turn_id错乱或顺序异常。 | 引入锁机制。例如,使用 Redis 分布式锁,键为f"session_lock:{session_id}"。在process_turn开始时加锁,结束时释放。 |
| 重建的 Session 无法继续执行,报“工具未找到” | 1. 持久化时,Step 的content中包含了工具类的实例,而非法序列化的配置。2. 重建后,工具注册表( ToolRegistry)与之前不同,缺少某些工具。 | 1. 检查持久化到数据库的 Stepcontent数据,是否是可序列化的字典。2. 对比服务重启前后的工具注册表。 | 1.只持久化配置:Step 的content应只存储工具名和参数字典,如{"tool": "get_weather", "parameters": {"city": "Beijing"}}。2.工具注册表应保持稳定:工具的实现和注册应在应用启动时完成,并保持一致。 |
6. 最佳实践与工程建议
基于 DeepSeek Harness 的设计和我们的实践,可以总结出以下构建生产级 Agent 系统的建议:
- 状态设计要幂等:Step 的执行应尽可能设计成幂等的。这样,在 Session 重建后,即使重复执行某个成功的 Step(在极端情况下),也不会产生副作用。例如,“发送邮件”这个 Action 应该先检查邮件是否已发送。
- 上下文管理策略:
Session.context不宜无限增长。对于长对话,需要设计上下文窗口管理策略,例如只保留最近 N 个 Turns 的摘要,或者将早期历史转移到长期记忆存储中。 - 持久化存储选型:
- 关系型数据库(如 PostgreSQL):结构清晰,利于复杂查询(如“查找所有失败的 Step”)。需要处理好三张表的关联查询和写入性能。
- 文档数据库(如 MongoDB):天然适合存储嵌套的 Session 文档,写入和重建快。但复杂查询可能需要索引支持。
- 时序数据库或专用存储:如果 Step 数据量极大,且需要分析执行链路性能,可考虑专用方案。
- 监控与可观测性:在每个 Step 中记录开始时间、结束时间和耗时。这不仅能用于性能监控,当 Session 重建后继续执行时,也有助于分析中断点。考虑集成 OpenTelemetry 等标准。
- 定义清晰的错误处理边界:区分可恢复错误(如网络超时)和不可恢复错误(如工具逻辑错误)。对于可恢复错误,可以在 Step 级别实现自动重试机制。
- 前端与状态同步:前端应用需要妥善处理
session_id。在单页应用(SPA)中,刷新页面后应能通过session_id从后端重建会话状态。WebSocket 连接断开重连时,也应携带session_id以恢复会话。
DeepSeek Harness 通过Turn/Step/Session的三层抽象,为 AI Agent 的复杂状态管理提供了一个清晰的蓝图。它本质上是一个面向状态的编程模型,将 Agent 的“思考流”物化为一系列可持久化、可查询、可回滚的步骤记录。
对于开发者而言,理解这一模型的价值远大于学会使用某个特定框架。它迫使你思考:你的 Agent 的“状态”是什么?如何定义它的“一步”?如何让它在任何意外中断后都能“无缝续杯”?当你开始用Session、Turn、Step的视角去设计系统时,你就已经跳出了简单的“请求-响应”模式,正在构建一个真正具有持续交互和记忆能力的智能体。
本文实现的简易系统是一个起点。你可以在此基础上,引入真实的 LLM 调用、复杂的工具集、更智能的步骤规划器(Planner)以及可视化调试界面。最终,你会拥有一个完全可控、可调试、可运维的 Agent 基础设施,这才是 AI 应用走向成熟的关键一步。