DeepAgents - Human in the loop

[复制链接]
发表于 2026-6-14 22:08:36 | 显示全部楼层 |阅读模式

马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。

您需要 登录 才可以下载或查看,没有账号?立即注册

×
前言

Human in the loop(人机协作)在企业级 Agent 应用中非常告急——AI 在执关键工具时必须颠末人类审批,克制误操纵影响业务。我之前用 LangGraph 0.3 裸写了一套(旧文),当时须要在 tool 函数里手动调 interrupt(),很啰嗦。如今有了 DeepAgents 和内置的 HumanInTheLoopMiddleware,只需设置一个 interrupt_on 字典,制止逻辑全主动——在实行前停息图实行,生存状态到 checkpointer,等待人类决定后规复。
不外官方文档的示例代码比力简朴,只演示了根本用法,没有阐明怎样在真实应用中整合。本文以"答应 AI 实行 shell 下令,但每次实行前需用户确认"为需求,一步步实现完备的人机协作流程。
焦点概念

在开始之前,先理清几个关键概念:
概念阐明interrupt(制止)当 Agent 预备调用某个被监控监控的 tool 时,HumanInTheLoopMiddleware 调用 LangGraph 的 interrupt() 停息图实行,并抛出包罗 action_requests 和 review_configs 的哀求checkpoint(查抄点)制止时图状态会被长期化。必须设置 checkpointer,否则制止后无法规复。生产环境发起用 AsyncPostgresSaver,测试用 InMemorySaverversion="v2"LangGraph 1.0 的 v2 模式,ainvoke() 返回 GraphOutput 对象(含 .interrupts 属性),astream() 的 updates 流中会出现 __interrupt__ 变乱Command(resume=)用户做出决定后,用 Command(resume={"decisions": [...]}) 从断点规复实行Decision(决定)四种范例:approve(答应)、reject(拒绝并反馈)、edit(修改参数后实行)、respond(人类直接回复,跳过 tool 实行)实行生命周期
  1. 用户提问 → Agent 调用 LLM 生成回复
  2.   → LLM 决定调用 tool(如 execute_shell_command)
  3.     → after_model 钩子:检查 tool 是否在 interrupt_on 中
  4.       → 是:构建 HITLRequest → interrupt() → 暂停 ⌛
  5.       → 否:继续执行
  6.   → 人类做出决策(approve / reject / edit / respond)
  7.     → 恢复执行 → 执行/拒绝 tool → LLM 生成最终回复 → 返回
复制代码
流程逻辑

以 Chainlit 谈天应用为交互载体,消息处理惩罚流程如下:

  • 用户在谈天页面发送消息(如"查抄下体系负载")
  • Agent 调用 LLM 天生复兴,LLM 决定调用 execute_shell_command
  • HumanInTheLoopMiddleware 检测到该 tool 在 interrupt_on 列表中,触发制止
  • Chainlit 应用检测到制止,向用户展示审批提示
  • 用户复兴 答应 / 拒绝,应用用 Command(resume=) 规复实行
  • Agent 根据决定实行或拒绝 tool,终极返回结果给用户
设置制止

起首须要在创建 Agent 时设置 HumanInTheLoopMiddleware:
  1. from deepagents import create_deep_agent
  2. from langchain.agents.middleware import HumanInTheLoopMiddleware
  3. agent = create_deep_agent(
  4.     model=llm,
  5.     tools=[execute_shell_command],
  6.     checkpointer=checkpointer,  # 必须配置!
  7.     system_prompt="你是一位智能助手...",
  8.     middleware=[
  9.         HumanInTheLoopMiddleware(
  10.             interrupt_on={
  11.                 # 对 execute_shell_command 进行审批
  12.                 "execute_shell_command": {
  13.                     "allowed_decisions": ["approve", "reject"]
  14.                 }
  15.             }
  16.         ),
  17.     ],
  18. )
复制代码
interrupt_on 是一个字典,key 为 tool 名称,value 的可设置项:

  • True — 答应全部四种决定(approve / edit / reject / respond)
  • False — 不拦截该 tool(等同于不写)
  • {"allowed_decisions": [...]} — 只答应指定决定范例
  • 还可以设置 when 谓词按参数条件判定是否拦截、description 自界说制止提示文本
invoke 模式中的实现

v2 模式下的 ainvoke() 返回 GraphOutput 对象,可通过 .interrupts 属性直接获取制止数据,不须要去查 state。
检测制止
  1. resp = await agent.ainvoke(
  2.     input={"messages": [HumanMessage(content=query)]},
  3.     config=config,
  4.     version="v2",
  5. )
  6. if resp.interrupts:
  7.     # 存在中断,resp.interrupts 是 Interrupt 对象的元组
  8.     interrupt = resp.interrupts[0]
  9.     # interrupt.value 是 HITLRequest,包含 action_requests 和 review_configs
  10.     print(interrupt.value["action_requests"])
复制代码
规复制止

用户做出决定后,用 Command(resume=) 规复:
  1. from langgraph.types import Command
  2. await agent.ainvoke(
  3.     Command(resume={
  4.         "decisions": [{"type": "approve"}]  # 或 {"type": "reject", "message": "..."}
  5.     }),
  6.     config=config,  # 必须用同一个 thread_id
  7.     version="v2",
  8. )
复制代码
关键:怎样分辨"新消息"照旧"制止规复"

在谈天应用中,用户发来的每条消息都走同一个 @cl.on_message 处理惩罚函数。用户说"查抄负载"和复兴"答应"都只是文本。办理方法是——调用前先查抄是否有待处理惩罚的制止:
  1. # 检查当前会话是否有待处理的中断
  2. state = await agent.aget_state(config)
  3. if state.next:
  4.     # 有待处理中断 → 本次消息是审批回复,构建 resume 命令
  5.     cmd = Command(resume={"decisions": [{"type": "approve"}]})
  6.     await agent.ainvoke(cmd, config=config, version="v2")
  7. else:
  8.     # 无中断 → 正常对话
  9.     resp = await agent.ainvoke(
  10.         {"messages": [HumanMessage(content=query)]}, config=config, version="v2"
  11.     )
复制代码
state.next 不为空表现图实行被停息了(有制止等待处理惩罚)。
stream 模式中的实现

流式模式须要用 stream_mode=["messages", "updates"](官方保举同时开启两种流):

  • messages 流:获取 LLM 的 token 级输出
  • updates 流:检测制止变乱 __interrupt__
  1. async for chunk in agent.astream(
  2.     input=input_data,
  3.     stream_mode=["messages", "updates"],
  4.     version="v2",
  5.     config=config,
  6. ):
  7.     if chunk["type"] == "messages":
  8.         msg, _meta = chunk["data"]
  9.         # msg 是 AIMessageChunk,包含 content 和 tool_calls
  10.         if isinstance(msg, AIMessageChunk) and msg.content:
  11.             yield extract_text(msg)     # 流式输出文本
  12.     elif chunk["type"] == "updates":
  13.         if "__interrupt__" in chunk["data"]:
  14.             interrupt = chunk["data"]["__interrupt__"][0]
  15.             yield format_question(interrupt)  # 输出审批问题
复制代码
stream 模式的规复与 invoke 类似——在调用 astream() 之前同样要先查抄 state.next 来判定是正常对话照旧制止规复。
完备示例

下面是焦点代码。checkpointer 和 llm 的设置函数、日志日志模块等非焦点代码省略。
Agent 封装(internal/agent/agent.py 焦点部分)
  1. # --- 审批关键词匹配 ---
  2. _APPROVE_KEYWORDS = frozenset(
  3.     {"yes", "accept", "approve", "ok", "是", "允许", "同意", "批准"}
  4. )
  5. def _parse_decision(query: str) -> str:
  6.     return "approve" if query.strip().lower() in _APPROVE_KEYWORDS else "reject"
  7. def _build_resume_command(decision_type: str, actions_count: int) -> Command:
  8.     item = {"type": decision_type}
  9.     if decision_type == "reject":
  10.         item["message"] = "user rejected this action"
  11.     return Command(resume={"decisions": [item for _ in range(actions_count)]})
  12. def _extract_text(message) -> str:
  13.     """从消息中提取纯文本(兼容 str 和 list[dict] 两种 content 格式)。"""
  14.     if not message or not hasattr(message, "content"):
  15.         return ""
  16.     content = message.content
  17.     if isinstance(content, str):
  18.         return content
  19.     if isinstance(content, list):
  20.         return "".join(
  21.             b.get("text", "")
  22.             for b in content
  23.             if isinstance(b, dict) and b.get("type") == "text"
  24.         )
  25.     return ""
  26. def _format_interrupt_question(interrupt) -> str:
  27.     """将中断数据格式化为用户的审批问题。"""
  28.     action_requests = interrupt.value.get("action_requests", [])
  29.     review_configs = interrupt.value.get("review_configs", [])
  30.     allowed = (
  31.         review_configs[0].get("allowed_decisions", ["approve", "reject"])
  32.         if review_configs
  33.         else ["approve", "reject"]
  34.     )
  35.     lines = []
  36.     for req in action_requests:
  37.         lines.append(
  38.             "Do you approve me to execute this action?\n\n"
  39.             f"- name: {req['name']}\n"
  40.             f"- args: `{req['args']}`\n"
  41.         )
  42.     lines.append(f"Input your decision: {', '.join(allowed)}\n")
  43.     return "\n".join(lines)
  44. class AIAgent:
  45.     # ... __init__, _init_deep_agent, _init_tools 省略 ...
  46.     async def _has_pending_interrupt(self, config: RunnableConfig) -> bool:
  47.         state = await self._agent.aget_state(config)
  48.         return bool(state.next)
  49.     # --- invoke 模式 ---
  50.     async def ainvoke(self, query: str, config: RunnableConfig) -> str:
  51.         if not self._agent:
  52.             await self._init_deep_agent()
  53.         # 优先处理中断恢复
  54.         if await self._has_pending_interrupt(config):
  55.             state = await self._agent.aget_state(config)
  56.             actions_count = len(
  57.                 state.interrupts[0].value["action_requests"]
  58.             )
  59.             decision = _parse_decision(query)
  60.             cmd = _build_resume_command(decision, actions_count)
  61.             await self._agent.ainvoke(cmd, config=config, version="v2")
  62.             # 恢复后取最新消息
  63.             state = await self._agent.aget_state(config)
  64.             if state.values and "messages" in state.values:
  65.                 return _extract_text(state.values["messages"][-1])
  66.             return "Oops, something went wrong."
  67.         # 正常对话
  68.         resp = await self._agent.ainvoke(
  69.             input={"messages": [HumanMessage(content=query)]},
  70.             config=config,
  71.             version="v2",
  72.         )
  73.         if resp.interrupts:
  74.             return _format_interrupt_question(resp.interrupts[0])
  75.         return _extract_text(resp.value["messages"][-1])
  76.     # --- stream 模式 ---
  77.     async def astream(self, query: str, config: RunnableConfig):
  78.         if not self._agent:
  79.             await self._init_deep_agent()
  80.         # 判断是中断恢复还是正常对话
  81.         state = await self._agent.aget_state(config)
  82.         if state.next:
  83.             actions_count = len(
  84.                 state.interrupts[0].value["action_requests"]
  85.             )
  86.             decision = _parse_decision(query)
  87.             input_data = _build_resume_command(decision, actions_count)
  88.         else:
  89.             input_data = {"messages": [HumanMessage(content=query)]}
  90.         async for chunk in self._agent.astream(
  91.             input=input_data,
  92.             stream_mode=["messages", "updates"],
  93.             version="v2",
  94.             config=config,
  95.         ):
  96.             if chunk["type"] == "messages":
  97.                 msg, _meta = chunk["data"]
  98.                 if isinstance(msg, AIMessageChunk) and msg.content:
  99.                     yield _extract_text(msg)
  100.             elif chunk["type"] == "updates" and "__interrupt__" in chunk["data"]:
  101.                 yield _format_interrupt_question(
  102.                     chunk["data"]["__interrupt__"][0]
  103.                 )
复制代码
Chainlit 应用层(chainlit_app.py 焦点部分)
  1. @cl.on_message
  2. async def main(msg: cl.Message):
  3.     config = RunnableConfig(
  4.         configurable={"thread_id": cl.context.session.id},
  5.     )
  6.     # stream 模式(推荐)
  7.     final_answer = cl.Message(content="")
  8.     async for chunk in ai_agent.astream(msg.content, config=config):
  9.         await final_answer.stream_token(chunk)
  10.     await final_answer.send()
  11.     # 或者 invoke 模式
  12.     # resp = await ai_agent.ainvoke(msg.content, config)
  13.     # await cl.Message(resp).send()
复制代码
Chainlit 应用层的代码非常简便——由于制止检测和规复逻辑全部封装在 AIAgent 内部了。Chainlit 只须要流式输出 astream() / ainvoke() 的返回结果即可。
交互结果
  1. [用户]: 检查下系统负载
  2. [AI]: 🔧 正在调用工具: execute_shell_command...
  3. [AI]: Do you approve me to execute this action?
  4.        - name: execute_shell_command
  5.        - args: `{"command": "cat /proc/loadavg && free -h", "timeout": 10}`
  6.        Input your decision: approve, reject
  7. [用户]: 批准
  8. [AI]: 当前系统负载: 0.52 0.38 0.25,
  9.        内存总容量 62Gi,已用 10Gi,剩余 46Gi,系统运行正常。
复制代码
进阶:利用 interrupt_on 的 when 谓词

假如不想拦截全部 shell 下令,只想拦截伤害操纵(如 rm、dd、写入体系目次等),可以用 when 谓词按参数条件判定:
  1. from langgraph.prebuilt.tool_node import ToolCallRequest
  2. def is_dangerous_command(request: ToolCallRequest) -> bool:
  3.     """只拦截包含危险操作的命令。"""
  4.     command = request.tool_call["args"].get("command", "")
  5.     dangerous = {"rm ", "dd ", "mkfs", "shutdown", "reboot"}
  6.     return any(d in command for d in dangerous)
  7. HumanInTheLoopMiddleware(
  8.     interrupt_on={
  9.         "execute_shell_command": {
  10.             "allowed_decisions": ["approve", "reject"],
  11.             "when": is_dangerous_command,   # 只在危险命令时拦截
  12.         }
  13.     }
  14. )
复制代码
when 谓词返回 True 才触发制止,返回 False 则主动答应。留意 when 须要 langchain >= 1.3.3。
改进点


  • 如今 reject 时用的是固定消息,现实产物中可以让用户输入拒绝缘故起因,方便 LLM 调解后续举动
  • 审批提示如今是纯文本,可以用 Chainlit 的 AskActionMessage 做成按钮交互(不外受制于 Chainlit Action 的 payload 范例限定,须要额外处理惩罚)
  • 假如有多个 tool 同时被拦截,action_requests 列表中会有多项,本文为简化只取了第一个,生产环境应遍历处理惩罚
增补

DeepAgents 的 HumanInTheLoopMiddleware 把之前须要手写的制止逻辑全部封装好了。集成到真实应用的关键只有三步:

  • 创建 Agent 时:设置 interrupt_on 字典 + 确保有 checkpointer
  • 每次调用前:通过 state.next 判定是正常对话照旧制止规复
  • 规复时:用 Command(resume={"decisions": [...]}) 传入用户决定
ainvoke 和 astream 两种模式的焦点逻辑同等,只是检测制止的方式差别(.interrupts 属性 vs updates 流中的 __interrupt__)。
回复

使用道具 举报

登录后关闭弹窗

登录参与点评抽奖  加入IT实名职场社区
去登录
快速回复 返回顶部 返回列表