
AutoGen 异步 Human-in-the-Loop 实战基于状态持久化实现“等待真实用户再执行工具调用”的 Agent 系统【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogen本文基于当前仓库 AutoGen.NET/Python agentic AI 框架中的异步人工介入示例 core_async_human_in_the_loop 展开讲解如何用autogen-core的RoutedAgent、SingleThreadedAgentRuntime与干预处理器Intervention Handler构建一个预约会议的双 Agent 系统——AI 助手先向慢速用户代理发出提问运行时挂起并将完整状态持久化待真实人类通过终端回复后重新水合rehydrate系统状态并继续执行工具调用直至任务终止。读完本文你将掌握异步 HITL 的完整代码骨架、状态保存/恢复的底层原理以及如何把该模式推广到邮件、审批流等更复杂的慢外部系统场景。1. 示例概览与设计动机关联文档core_async_human_in_the_loop/README.md核心代码main.py模型配置模板model_config_template.yml。该示例解决一个真实工程痛点人在回路中的系统是慢的。人类响应需要时间且不同沟通媒介终端、邮件、IM、审批流的延迟差异极大。当 Agent 等待人类回复期间进程可能被回收、容器可能被重启、应用可能被关闭。此时必须把系统运行状态连同正在等待什么的上下文一并持久化当用户回来时再按需重建状态继续运行。该示例的核心诉求在进行工具调用之前暂停流程先等待人类输入确认或补充缺失参数将运行时runtime中所有已实例化 Agent 的状态保存到持久层收到人类回复后用保存的状态重新水合运行时并继续执行。整个系统只有两个 Agent一个负责调度、可调用schedule_meeting工具的 AI 助手一个充当慢速人类代理的 User Proxy。同时主文档与 main.py 的模块说明明确指出本示例把人抽象为慢系统同样的原理可直接应用于任何 Agent 需要交互的慢速外部系统。2. 运行环境与安装示例基于 Python 生态的 AutoGen 核心包。关联文档给出了明确的安装前提——需要一个已安装autogen-core及所需依赖的 shell 环境pip install autogen-ext[openai,azure] pyyamlautogen-ext[openai,azure]提供 OpenAI / Azure OpenAI 的模型客户端扩展组件示例中的模型客户端即来自该扩展pyyaml用于读取model_config.yml配置文件。运行也非常简单见 README.mdpython main.pymain.py的__main__入口从当前目录读取model_config.yml加载模型配置因此运行前必须准备配置文件详见下一节。3. 模型配置从模板生成 model_config.yml示例运行前需要一份模型配置。仓库提供了模板 model_config_template.yml其中内置了三种可切换的接入方式方式一OpenAI 直接使用 API Keyprovider: autogen_ext.models.openai.OpenAIChatCompletionClient config: model: gpt-4o api_key: REPLACE_WITH_YOUR_API_KEY方式二Azure OpenAI 使用 API Keyprovider: autogen_ext.models.openai.AzureOpenAIChatCompletionClient config: model: gpt-4o azure_endpoint: https://{your-custom-endpoint}.openai.azure.com/ azure_deployment: {your-azure-deployment} api_version: {your-api-version} api_key: REPLACE_WITH_YOUR_API_KEY方式三Azure OpenAI 使用 Azure AD Token Providerprovider: autogen_ext.models.openai.AzureOpenAIChatCompletionClient config: model: gpt-4o azure_endpoint: https://{your-custom-endpoint}.openai.azure.com/ azure_deployment: {your-azure-deployment} api_version: {your-api-version} azure_ad_token_provider: provider: autogen_ext.auth.azure.AzureTokenProvider config: provider_kind: DefaultAzureCredential scopes: - https://cognitiveservices.azure.com/.default实际使用时应把模板复制为同目录下的model_config.yml替换其中的占位符REPLACE_WITH_YOUR_API_KEY、endpoint、deployment、api_version 等。在 main.py 中通过yaml.safe_load读取该文件并交给ChatCompletionClient.load_component(model_config)完成组件化加载见 main.py——这是 AutoGen 组件Component化配置体系的典型用法provider指向类路径config中的字段与组件配置模型一一对应从而无需在业务代码中写死任何供应商细节。4. 整体架构与两种运行分支系统完整流程由 main() 驱动。核心组件包括组件角色源码位置SlowUserProxyAgent订阅同一 topic充当慢速人类的代理把助手问题转交给真实用户main.pySchedulingAssistantAgentAI 调度助手持有模型客户端与工具决定是调用工具还是发问main.pyNeedsUserInputHandler拦截GetSlowUserMessage标记需要用户输入main.pyTerminationHandler拦截TerminateMessage标记任务已终止main.pyMockPersistence内存版持久层可替换为文件/数据库main.pySingleThreadedAgentRuntime单线程事件循环运行时注册干预处理器main.py这两个 Agent 通过type_subscription(scheduling_assistant_conversation)同时订阅同一个主题因此双方都可以通过publish_message(...)把消息广播到DefaultTopicId(scheduling_assistant_conversation)上被对方接收——这是 AutoGen Core 的主题订阅 发布订阅通信模型type_subscription装饰器定义见 _default_subscription.py。运行时存在两条分支首轮启动无用户输入主程序向调度助手投递一条AssistantTextMessageHi! How can I help you? I can help schedule meetings作为会话种子消息助手读取后开始对话恢复继续带用户输入说明在上一轮运行时被挂起处有等待回答的问题程序先runtime.load_state(state)恢复此前保存的全部 Agent 状态再把人类的新输入封装为UserTextMessage投递给调度助手让它带着完整记忆继续运行。从源码注释main.py可以看到这是一套刻意简化、却保留了完整关键点的骨架可平滑扩展为更复杂的场景。5. 消息类型设计把会话与控制信号显式建模示例将对话内容与流程控制信号都建模为 dataclass 消息这是事件驱动 Agent 系统可读性的关键dataclass class TextMessage: source: str content: str dataclass class UserTextMessage(TextMessage): pass dataclass class AssistantTextMessage(TextMessage): pass dataclass class GetSlowUserMessage: content: str # 真正要抛给人类的问题 dataclass class TerminateMessage: content: str # 任务完成信号设计要点UserTextMessage/AssistantTextMessage分别代表进入系统的真实人类输入和助手产出的文本二者区分了消息的方向与来源GetSlowUserMessage不投递给任何 Agent而是作为干预处理器捕获的信号——当它在发布publish_message环节流经干预处理器时被识别并记录从而让运行时的宿主代码知道该停下来问人了TerminateMessage同理作为任务成功完成的终止信号。6. 慢速用户代理 SlowUserProxyAgent 的实现type_subscription(scheduling_assistant_conversation) class SlowUserProxyAgent(RoutedAgent): def __init__(self, name: str, description: str) - None: super().__init__(description) self._model_context BufferedChatCompletionContext(buffer_size5) self._name name message_handler async def handle_message(self, message: AssistantTextMessage, ctx: MessageContext) - None: await self._model_context.add_message(AssistantMessage(contentmessage.content, sourcemessage.source)) await self.publish_message( GetSlowUserMessage(contentmessage.content), topic_idDefaultTopicId(scheduling_assistant_conversation), ) async def save_state(self) - Mapping[str, Any]: return {memory: await self._model_context.save_state()} async def load_state(self, state: Mapping[str, Any]) - None: await self._model_context.load_state(state[memory])它做了三件事保存收到的助手消息到本地会话上下文——每个 Agent 各自维护一份BufferedChatCompletionContext立即把问题以GetSlowUserMessage形式重新发布到同一 topic。该消息没有对应的 Agent 消息处理器因此唯一会接住它的只有干预处理器从而把向人类提问从业务逻辑中解耦出来实现save_state/load_state状态协议将上下文序列化内容交给运行时统一管理。这里体现的抽象是SlowUserProxyAgent本身不真正向人类询问——询问动作被发布一条无处理器消息这一事件表达出来再由宿主代码决定以何种介质终端、邮件、API 等触达真实人类。7. 调度助手 SchedulingAssistantAgent工具调用前的等待点type_subscription(scheduling_assistant_conversation) class SchedulingAssistantAgent(RoutedAgent): def __init__(self, name, description, model_client, initial_messageNone): ... self._model_context BufferedChatCompletionContext( buffer_size5, initial_messages[UserMessage(...)] if initial_message else None, ) self._system_messages [SystemMessage(contentf...今天日期是 {datetime.datetime.now().strftime(%Y-%m-%d)})]其消息处理逻辑main.py分两种情况情况 A模型请求调用工具返回 FunctionCall 列表response await self._model_client.create( self._system_messages (await self._model_context.get_messages()), toolstools ) if isinstance(response.content, list) and all(isinstance(item, FunctionCall) for item in response.content): for call in response.content: tool next((tool for tool in tools if tool.name call.name), None) if tool is None: raise ValueError(fTool not found: {call.name}) arguments json.loads(call.arguments) await tool.run_json(arguments, ctx.cancellation_token, call_idcall.id) await self.publish_message( TerminateMessage(contentMeeting scheduled), topic_idDefaultTopicId(scheduling_assistant_conversation), ) return即助手向模型提供tools[ScheduleMeetingTool()]若模型决定调用工具直接解析参数并执行run_json成功后发布TerminateMessage宣告完成。情况 B模型仍需要更多信息返回纯文本此时助手把文本作为AssistantTextMessage发布并写入自己的会话上下文。随后该消息被SlowUserProxyAgent接收进而转换为对真实人类的提问——流程正是在这里挂起等待人类回答。值得注意的细节ScheduleMeetingTool是BaseTool[ScheduleMeetingInput, ScheduleMeetingOutput]的子类输入输出均以 Pydantic 模型声明recipient、date、time三个带描述字段run方法接收强类型的ScheduleMeetingInputmain.py。这展示了 AutoGen Core 工具系统结构化参数 → 自动 JSON 解析 → 类型安全调用的链路。8. 干预处理器让运行时感知外部世界AutoGen Core 的干预处理器允许在消息经运行时提交时对其进行修改、记录或丢弃协议定义见 _intervention.py。InterventionHandler定义了三类可拦截点on_send发送、on_publish发布、on_response响应。DefaultInterventionHandler提供了默认空实现便于子类只覆写关心的方法见 _intervention.py。示例中NeedsUserInputHandler与TerminationHandler都覆写了on_publishclass NeedsUserInputHandler(DefaultInterventionHandler): async def on_publish(self, message, *, message_context): if isinstance(message, GetSlowUserMessage): self.question_for_user message # 捕获待转交人类的提问 return message property def needs_user_input(self) - bool: return self.question_for_user is not None property def user_input_content(self) - str | None: return self.question_for_user.content if self.question_for_user else NoneTerminationHandler结构完全对称只是捕获对象换成TerminateMessage并暴露is_terminated与termination_msg。这两个处理器在创建运行时时就注入main.pyruntime SingleThreadedAgentRuntime(intervention_handlers[needs_user_input_handler, termination_handler])需要注意处理器方法返回None会被视为无变更并产生警告若要真正丢弃消息必须显式返回DropMessage见 _intervention.py。当前仓库中明确支持干预处理器的运行时是SingleThreadedAgentRuntime单线程运行时文档注释它适合单机开发与独立应用场景。9. 主循环状态机、挂起与恢复main()是贯穿整个流程的状态机控制器main.py核心步骤为加载模型客户端model_client ChatCompletionClient.load_component(model_config)实例化两个干预处理器与运行时并注册两个 AgentUser与SchedulingAssistant后者携带初始问候消息作为会话种子判定运行分支latest_user_input为空 → 投递初始AssistantTextMessage开启会话非空 → 构造UserTextMessage恢复状态如果存在if state: await runtime.load_state(state)——把上一轮保存的各 Agent 内存重新装回发布启动消息并runtime.start()阻塞等待停止条件await runtime.stop_when(lambda: termination_handler.is_terminated or needs_user_input_handler.needs_user_input)即要么任务完成要么需要人类输入二者居一时停止消息循环 7.保存状态并返回待答问题根据两个处理器的状态判定结果若needs_user_input_handler.user_input_content非空则以返回值形式把问题交给外层循环否则打印终止信息。无论哪种情况都执行state_to_persist await runtime.save_state()并写入持久层保证下一轮可恢复。在__main__中则用递归方式串联多轮提问-回答main.pyasync def run_main(question_for_user: str | None None): if question_for_user: user_input get_user_input(question_for_user) # 终端阻塞输入 else: user_input None user_input_needed await main(model_config, user_input) if user_input_needed: await run_main(user_input_needed) # 带着新回答再次进入 runtime asyncio.run(run_main())终端交互输出大致为--------------------------QUESTION_FOR_USER-------------------------- Hi! I can help schedule meetings. Please provide the recipient and date. --------------------------------------------------------------------- Enter your input: ......get_user_input用ThreadPoolExecutor包装标准input()为异步等待ainput 的等价模式既保持进程阻塞等待人类又不干扰 asyncio 事件循环。10. 状态保存与恢复的运行时级原理示例把可重入性建立在运行时的状态协议上。查看 SingleThreadedAgentRuntime 的实现save_state()L431-L447会遍历所有已实例化的 Agent逐一调用其save_state()最终返回以 Agent ID 为键、各 Agent 状态为值的字典load_state()L449-L464遍历状态字典对每个AgentId.from_str解析出的、且当前运行时已知的 Agent 类型调用load_state(state)回填需要指出的是源码注释明确说明当前版本不会保存订阅状态subscription state计划在未来补充见 L437-L439。因此若 Agent 的订阅是静态声明的本示例即如此保存/恢复不受影响若订阅是动态变化的重水合后需注意这一点。单个 Agent 侧的状态则落在BufferedChatCompletionContext其保存/恢复序列化的是最近 N 条消息的内存窗口。buffer_size5意味着每次送入模型的历史被裁剪为最近 5 条其get_messages()返回self._messages[-self._buffer_size:]见 _buffered_chat_completion_context.py。而示例中的MockPersistence仅是一个dict包装的内存持久层用于演示持久化接口load_content/save_content在生产中可将其替换为文件、数据库、对象存储等真正的持久化实现或者存储到远端以便跨进程/跨主机恢复。10.1 关于 stop_when 的一个工程提示从运行时源码可见stop_when是通过轮询条件实现的其文档注释明确提示该方法不推荐用于新代码保留它出于历史原因它会以 busy loop 持续检查条件建议使用stop_when_idle或stop若需按条件停止更高效的做法是用后台任务与asyncio.Event发出信号见 L856-L868。示例中仍使用stop_when(lambda: ...)来同时监听终止与需要输入两个条件这是为保持教学骨架简洁而做的取舍。读者在把它改造成生产代码时可参考上述建议改用asyncio.Event驱动停止以消除忙等开销。11. 设计权衡与扩展思路结合示例源码顶部注释main.py与代码实现可以提炼出如下可迁移的设计原则决定在哪一个时点持久化是关键权衡。本示例选择在需要用户输入的瞬间保存状态——因为这是唯一确定的挂起点而在更复杂的系统中可能需要在多个关键节点保存以保证系统能被回滚/恢复到正确状态。保存点越多恢复越精细但 I/O 与复杂度也越高把等待人泛化为等待慢系统。GetSlowUserMessage→ 干预处理器 → 持久化 → 外部触达 这套链路上外部触达可以是终端输入当前实现、也可以是发邮件、推送 IM、写入审批队列等任何慢通道回复内容再作为latest_user_input喂回main()。因此同一套骨架即可支撑异步审批、人类复核工具调用、跨平台多模态人工介入等场景每个 Agent 自持上下文并实现状态协议save_state/load_state使运行时得以无差别地统一调度这是把 HITL 做得可重入、可分布的前提干预处理器让宿主代码与Agent 业务解耦Agent 只负责发布语义化信号需要输入 / 任务终止要不要停、何时停、如何触达人类全由宿主main()决定方便复用与测试。12. 如何进一步阅读与运行验证完整可运行代码main.py模型配置模板含 OpenAI / Azure Key / Azure AD Token 三种模式model_config_template.yml干预处理器协议与默认实现autogen_core/_intervention.py运行时保存/恢复与停止逻辑autogen_core/_single_threaded_agent_runtime.py有界会话上下文消息缓冲与裁剪autogen_core/model_context/_buffered_chat_completion_context.py主题订阅装饰器autogen_core/_default_subscription.py。自行复现时请务必把 model_config_template.yml 复制为同目录下model_config.yml并填入有效凭证随后安装autogen-ext[openai,azure]与pyyaml执行python main.py即可观察助手提问 → 系统挂起 → 人类回答 → 状态恢复 → 执行预约工具 → 正常终止的完整异步 HITL 流程。若需要开启调试日志示例中预留了logging与状态文件清理的注释代码取消注释即可main.py。【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogen创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考