马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。
您需要 登录 才可以下载或查看,没有账号?立即注册
×
前言
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 实行)实行生命周期
- 用户提问 → Agent 调用 LLM 生成回复
- → LLM 决定调用 tool(如 execute_shell_command)
- → after_model 钩子:检查 tool 是否在 interrupt_on 中
- → 是:构建 HITLRequest → interrupt() → 暂停 ⌛
- → 否:继续执行
- → 人类做出决策(approve / reject / edit / respond)
- → 恢复执行 → 执行/拒绝 tool → LLM 生成最终回复 → 返回
复制代码 流程逻辑
以 Chainlit 谈天应用为交互载体,消息处理惩罚流程如下:
- 用户在谈天页面发送消息(如"查抄下体系负载")
- Agent 调用 LLM 天生复兴,LLM 决定调用 execute_shell_command
- HumanInTheLoopMiddleware 检测到该 tool 在 interrupt_on 列表中,触发制止
- Chainlit 应用检测到制止,向用户展示审批提示
- 用户复兴 答应 / 拒绝,应用用 Command(resume=) 规复实行
- Agent 根据决定实行或拒绝 tool,终极返回结果给用户
设置制止
起首须要在创建 Agent 时设置 HumanInTheLoopMiddleware:- from deepagents import create_deep_agent
- from langchain.agents.middleware import HumanInTheLoopMiddleware
- agent = create_deep_agent(
- model=llm,
- tools=[execute_shell_command],
- checkpointer=checkpointer, # 必须配置!
- system_prompt="你是一位智能助手...",
- middleware=[
- HumanInTheLoopMiddleware(
- interrupt_on={
- # 对 execute_shell_command 进行审批
- "execute_shell_command": {
- "allowed_decisions": ["approve", "reject"]
- }
- }
- ),
- ],
- )
复制代码 interrupt_on 是一个字典,key 为 tool 名称,value 的可设置项:
- True — 答应全部四种决定(approve / edit / reject / respond)
- False — 不拦截该 tool(等同于不写)
- {"allowed_decisions": [...]} — 只答应指定决定范例
- 还可以设置 when 谓词按参数条件判定是否拦截、description 自界说制止提示文本
invoke 模式中的实现
v2 模式下的 ainvoke() 返回 GraphOutput 对象,可通过 .interrupts 属性直接获取制止数据,不须要去查 state。
检测制止
- resp = await agent.ainvoke(
- input={"messages": [HumanMessage(content=query)]},
- config=config,
- version="v2",
- )
- if resp.interrupts:
- # 存在中断,resp.interrupts 是 Interrupt 对象的元组
- interrupt = resp.interrupts[0]
- # interrupt.value 是 HITLRequest,包含 action_requests 和 review_configs
- print(interrupt.value["action_requests"])
复制代码 规复制止
用户做出决定后,用 Command(resume=) 规复:- from langgraph.types import Command
- await agent.ainvoke(
- Command(resume={
- "decisions": [{"type": "approve"}] # 或 {"type": "reject", "message": "..."}
- }),
- config=config, # 必须用同一个 thread_id
- version="v2",
- )
复制代码 关键:怎样分辨"新消息"照旧"制止规复"
在谈天应用中,用户发来的每条消息都走同一个 @cl.on_message 处理惩罚函数。用户说"查抄负载"和复兴"答应"都只是文本。办理方法是——调用前先查抄是否有待处理惩罚的制止:- # 检查当前会话是否有待处理的中断
- state = await agent.aget_state(config)
- if state.next:
- # 有待处理中断 → 本次消息是审批回复,构建 resume 命令
- cmd = Command(resume={"decisions": [{"type": "approve"}]})
- await agent.ainvoke(cmd, config=config, version="v2")
- else:
- # 无中断 → 正常对话
- resp = await agent.ainvoke(
- {"messages": [HumanMessage(content=query)]}, config=config, version="v2"
- )
复制代码 state.next 不为空表现图实行被停息了(有制止等待处理惩罚)。
stream 模式中的实现
流式模式须要用 stream_mode=["messages", "updates"](官方保举同时开启两种流):
- messages 流:获取 LLM 的 token 级输出
- updates 流:检测制止变乱 __interrupt__
- async for chunk in agent.astream(
- input=input_data,
- stream_mode=["messages", "updates"],
- version="v2",
- config=config,
- ):
- if chunk["type"] == "messages":
- msg, _meta = chunk["data"]
- # msg 是 AIMessageChunk,包含 content 和 tool_calls
- if isinstance(msg, AIMessageChunk) and msg.content:
- yield extract_text(msg) # 流式输出文本
- elif chunk["type"] == "updates":
- if "__interrupt__" in chunk["data"]:
- interrupt = chunk["data"]["__interrupt__"][0]
- yield format_question(interrupt) # 输出审批问题
复制代码 stream 模式的规复与 invoke 类似——在调用 astream() 之前同样要先查抄 state.next 来判定是正常对话照旧制止规复。
完备示例
下面是焦点代码。checkpointer 和 llm 的设置函数、日志 模块等非焦点代码省略。
Agent 封装(internal/agent/agent.py 焦点部分)
- # --- 审批关键词匹配 ---
- _APPROVE_KEYWORDS = frozenset(
- {"yes", "accept", "approve", "ok", "是", "允许", "同意", "批准"}
- )
- def _parse_decision(query: str) -> str:
- return "approve" if query.strip().lower() in _APPROVE_KEYWORDS else "reject"
- def _build_resume_command(decision_type: str, actions_count: int) -> Command:
- item = {"type": decision_type}
- if decision_type == "reject":
- item["message"] = "user rejected this action"
- return Command(resume={"decisions": [item for _ in range(actions_count)]})
- def _extract_text(message) -> str:
- """从消息中提取纯文本(兼容 str 和 list[dict] 两种 content 格式)。"""
- if not message or not hasattr(message, "content"):
- return ""
- content = message.content
- if isinstance(content, str):
- return content
- if isinstance(content, list):
- return "".join(
- b.get("text", "")
- for b in content
- if isinstance(b, dict) and b.get("type") == "text"
- )
- return ""
- def _format_interrupt_question(interrupt) -> str:
- """将中断数据格式化为用户的审批问题。"""
- action_requests = interrupt.value.get("action_requests", [])
- review_configs = interrupt.value.get("review_configs", [])
- allowed = (
- review_configs[0].get("allowed_decisions", ["approve", "reject"])
- if review_configs
- else ["approve", "reject"]
- )
- lines = []
- for req in action_requests:
- lines.append(
- "Do you approve me to execute this action?\n\n"
- f"- name: {req['name']}\n"
- f"- args: `{req['args']}`\n"
- )
- lines.append(f"Input your decision: {', '.join(allowed)}\n")
- return "\n".join(lines)
- class AIAgent:
- # ... __init__, _init_deep_agent, _init_tools 省略 ...
- async def _has_pending_interrupt(self, config: RunnableConfig) -> bool:
- state = await self._agent.aget_state(config)
- return bool(state.next)
- # --- invoke 模式 ---
- async def ainvoke(self, query: str, config: RunnableConfig) -> str:
- if not self._agent:
- await self._init_deep_agent()
- # 优先处理中断恢复
- if await self._has_pending_interrupt(config):
- state = await self._agent.aget_state(config)
- actions_count = len(
- state.interrupts[0].value["action_requests"]
- )
- decision = _parse_decision(query)
- cmd = _build_resume_command(decision, actions_count)
- await self._agent.ainvoke(cmd, config=config, version="v2")
- # 恢复后取最新消息
- state = await self._agent.aget_state(config)
- if state.values and "messages" in state.values:
- return _extract_text(state.values["messages"][-1])
- return "Oops, something went wrong."
- # 正常对话
- resp = await self._agent.ainvoke(
- input={"messages": [HumanMessage(content=query)]},
- config=config,
- version="v2",
- )
- if resp.interrupts:
- return _format_interrupt_question(resp.interrupts[0])
- return _extract_text(resp.value["messages"][-1])
- # --- stream 模式 ---
- async def astream(self, query: str, config: RunnableConfig):
- if not self._agent:
- await self._init_deep_agent()
- # 判断是中断恢复还是正常对话
- state = await self._agent.aget_state(config)
- if state.next:
- actions_count = len(
- state.interrupts[0].value["action_requests"]
- )
- decision = _parse_decision(query)
- input_data = _build_resume_command(decision, actions_count)
- else:
- input_data = {"messages": [HumanMessage(content=query)]}
- async for chunk in self._agent.astream(
- input=input_data,
- stream_mode=["messages", "updates"],
- version="v2",
- config=config,
- ):
- if chunk["type"] == "messages":
- msg, _meta = chunk["data"]
- if isinstance(msg, AIMessageChunk) and msg.content:
- yield _extract_text(msg)
- elif chunk["type"] == "updates" and "__interrupt__" in chunk["data"]:
- yield _format_interrupt_question(
- chunk["data"]["__interrupt__"][0]
- )
复制代码 Chainlit 应用层(chainlit_app.py 焦点部分)
- @cl.on_message
- async def main(msg: cl.Message):
- config = RunnableConfig(
- configurable={"thread_id": cl.context.session.id},
- )
- # stream 模式(推荐)
- final_answer = cl.Message(content="")
- async for chunk in ai_agent.astream(msg.content, config=config):
- await final_answer.stream_token(chunk)
- await final_answer.send()
- # 或者 invoke 模式
- # resp = await ai_agent.ainvoke(msg.content, config)
- # await cl.Message(resp).send()
复制代码 Chainlit 应用层的代码非常简便——由于制止检测和规复逻辑全部封装在 AIAgent 内部了。Chainlit 只须要流式输出 astream() / ainvoke() 的返回结果即可。
交互结果
- [用户]: 检查下系统负载
- [AI]: 🔧 正在调用工具: execute_shell_command...
- [AI]: Do you approve me to execute this action?
- - name: execute_shell_command
- - args: `{"command": "cat /proc/loadavg && free -h", "timeout": 10}`
- Input your decision: approve, reject
- [用户]: 批准
- [AI]: 当前系统负载: 0.52 0.38 0.25,
- 内存总容量 62Gi,已用 10Gi,剩余 46Gi,系统运行正常。
复制代码 进阶:利用 interrupt_on 的 when 谓词
假如不想拦截全部 shell 下令,只想拦截伤害操纵(如 rm、dd、写入体系目次等),可以用 when 谓词按参数条件判定:- from langgraph.prebuilt.tool_node import ToolCallRequest
- def is_dangerous_command(request: ToolCallRequest) -> bool:
- """只拦截包含危险操作的命令。"""
- command = request.tool_call["args"].get("command", "")
- dangerous = {"rm ", "dd ", "mkfs", "shutdown", "reboot"}
- return any(d in command for d in dangerous)
- HumanInTheLoopMiddleware(
- interrupt_on={
- "execute_shell_command": {
- "allowed_decisions": ["approve", "reject"],
- "when": is_dangerous_command, # 只在危险命令时拦截
- }
- }
- )
复制代码 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__)。
|