尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Agent Zero 消息队列发送端点解析:message_queue_send 的请求契约、聚合发送与队列生命周期

Agent Zero 消息队列发送端点解析:message_queue_send 的请求契约、聚合发送与队列生命周期 Agent Zero 消息队列发送端点解析message_queue_send 的请求契约、聚合发送与队列生命周期【免费下载链接】agent-zeroAgent Zero AI framework项目地址: https://gitcode.com/GitHub_Trending/ag/agent-zero本文围绕 Agent Zero 仓库中的api/message_queue_send.py.dox.md这份文件级 DOX 档案展开完整讲解/message_queue_send端点的请求参数、响应契约与四条执行分支并结合 api/message_queue_send.py 与队列核心实现 helpers/message_queue.py 深入剖析“单条发送”与“聚合批量发送”两种模式的底层原理。读完本文你可以准确理解 Agent Zero 是如何让用户在 Agent 忙时排队、闲时立即投递消息以及状态同步机制如何把队列变化推送到前端的。1. 消息队列在 Agent Zero 中的位置Agent Zero 的交互模型中当 Agent 正在处理一个任务monologue时用户新输入的消息不会直接打断它而是先进入消息队列message queue。队列围绕一次会话的AgentContext组织完整生命周期由三组 API 端点支撑端点实现文件职责POST /message_queue_addapi/message_queue_add.py把文本与附件加入队列返回item_id与当前队列长度POST /message_queue_sendapi/message_queue_send.py立即发送队列中指定的消息或聚合发送全部消息POST /message_queue_removeapi/message_queue_remove.py移除单条消息或不传item_id时清空整个队列队列数据并不存在独立数据库中而是挂在会话上下文的持久化数据区。从 helpers/message_queue.py 可以看到队列以QUEUE_KEY message_queue为键存储在context.data中另有一个自增序号键QUEUE_SEQ_KEY记录入队顺序附件文件名则相对固定的上传目录/a0/usr/uploads解析为完整路径见 helpers/message_queue.py。这意味着队列与某个 context 强绑定、跨重启可恢复。2. 端点实现与运行契约message_queue_send.py.dox.md作为该端点的“持久化档案”durable notes声明了它的责任边界、类结构与关键依赖运行实现由 api/message_queue_send.py 持有DOX 文件负责沉淀职责、契约、副作用与验证方式二者必须同步更新核心类MessageQueueSend继承自helpers.api.ApiHandler对外暴露async process(self, input: dict, request: Request) - dict | Response观察到的副作用领域是“settings/state persistence”——即每次成功发送都会触发状态标记驱动前端同步导入依赖面为agent、helpers、helpers.api、helpers.state_monitor_integration。端点完整实现非常精炼共 33 行api/message_queue_send.pyclass MessageQueueSend(ApiHandler): Send queued message(s) immediately. async def process(self, input: dict, request: Request) - dict | Response: context AgentContext.get(input.get(context, )) if not context: return Response(Context not found, status404) if not mq.has_queue(context): return {ok: True, message: Queue empty} item_id input.get(item_id) send_all input.get(send_all, False) if send_all: count mq.send_all_aggregated(context) if count: mark_dirty_for_context(context.id, reasonmessage_queue_send_all) return {ok: True, sent_count: count} # Send single item item mq.pop_item(context, item_id) if item_id else mq.pop_first(context) if not item: return Response(Item not found, status404) mq.send_message(context, item) mark_dirty_for_context(context.id, reasonmessage_queue_send) return {ok: True, sent_item_id: item[id]}2.1 请求参数参数类型必填说明contextstring是会话上下文 ID用于AgentContext.get()定位会话找不到返回 404 纯文本响应item_idstring否要立即发送的队列条目 ID仅在不传send_all时生效send_allbool否默认False。为True时忽略item_id把整队消息聚合成一条发出两个选填参数都不传时行为等价于“发送队首”item_id为None会走mq.pop_first(context)分支。2.2 响应契约与四条分支端点的控制流由三个前置检查加一个模式开关组成对应四种类别的结果上下文不存在AgentContext.get返回空直接返回Response(Context not found, status404)——这里用的是helpers.api.Response对象而非 JSON dict符合 DOX 中“非 JSON 响应用helpers.api.Response”的工作约定队列为空mq.has_queue(context)为假时返回{ok: true, message: Queue empty}注意这不是错误而是幂等友好的成功响应聚合发送send_allTrue时调用mq.send_all_aggregated(context)并返回{ok: true, sent_count: N}单条发送pop_item/pop_first取出的条目若为空并发下已被取走返回 404Item not found否则发送成功并返回{ok: true, sent_item_id: id}。DOX 同时强调HTTP handler 必须继承helpers.api.ApiHandlerWebSocket handler 则继承helpers.ws.WsHandler并且认证、CSRF、loopback 与 API-key 检查属于框架层契约除非端点契约明确变更否则必须原样保留。请求载荷、认证要求、响应形状或路由副作用一旦变化DOX 文件也要同步更新——这是仓库中“扁平目录 每实现文件一份.dox.md”约定的核心纪律。3. 两种发送模式的底层实现单条与批量两条路径最终都汇聚到 helpers/message_queue.py 中的两个核心函数。3.1 单条发送pop 后投递mq.send_messagehelpers/message_queue.py的语义是“记录 投递”def send_message(context, item, source (from queue)): from agent import UserMessage # 延迟导入避免循环依赖 message item.get(text, ) attachments item.get(attachments, []) msg_id str(uuid.uuid4()) log_user_message(context, message, attachments, message_idmsg_id, sourcesource) context.communicate(UserMessage(message, attachments, idmsg_id))log_user_message先在控制台以紫底白字样式打印再写入 UI 日志typeuser保证用户能看到消息“已投递”context.communicate(UserMessage(...))才是真正把消息注入 Agent 输入通道的动作UserMessage采用函数内延迟导入DOX 提到的依赖面中agent正是这条循环依赖的由来消息使用新的uuid4作为日志 ID与队列条目 ID 解耦。取条目的两个函数都遵循“弹出 落盘 同步前端”的固定三步pop_firsthelpers/message_queue.pyqueue.pop(0)弹出队首写回context.set_data再执行_sync_outputpop_itemhelpers/message_queue.py按 ID 线性查找并queue.pop(i)找不到返回None端点据此给出 404。3.2 聚合发送合并为一条消息send_all_aggregatedhelpers/message_queue.py展示了“批量不是连发 N 条而是合并成 1 条”的设计text \n\n---\n\n.join(i[text] for i in items if i[text]) attachments [a for i in items for a in i.get(attachments, [])] ... log_user_message(context, text, attachments, message_idmsg_id, source (queued batch)) context.communicate(UserMessage(text, attachments, idmsg_id)) return len(items)循环pop_first排空整个队列各条文本以\n\n---\n\n分隔线拼接空文本条目被跳过所有条目的附件按队列顺序拍平合并为一个列表返回条目数量端点据此填充sent_count。这种聚合策略让 Agent 只产生一次完整回复来响应整批排队消息而不是逐条应答显著减少上下文往返次数。3.3 队列状态如何同步到前端每次增删改都会触发_sync_outputhelpers/message_queue.py它把队列投影为前端友好的“截断视图”——文本超过 100 字符自动加...附件只保留文件名并附带attachment_count——写入context.set_output_data(QUEUE_KEY, truncated)。WebUI 端通过轮询拿到这份 output data 渲染队列面板这正是 DOX 所说“settings/state persistence”副作用的具体落点。4. 前端调用方message-queue-storeWebUI 侧的调用集中在 webui/components/chat/message-queue/message-queue-store.js与本端点直接相关的有两个方法async sendItem(itemId) { const context globalThis.getContext?.(); if (!context) return; await api.callJsonApi(/message_queue_send, { context, item_id: itemId }); } async sendAll() { // 存在上传中的 pending 条目时先提示用户等待或移除 await api.callJsonApi(/message_queue_send, { context, send_all: true }); }sendItem对应“单条立即发送”分支与 3.1 的pop_item路径一一映射sendAll在调用前先检查本地pendingItems上传尚未完成、队列中还是临时 ID 的条目若有则弹出 toast 提示“等待上传完成或移除”这正是端点参数send_all的前端生产者入队侧addToQueue会把前端生成的临时tempId作为item_id传给/message_queue_add见 api/message_queue_add.py 中mq.add(context, text, attachments, item_id)的透传使得“pending 条目”与“服务端条目”可以凭同一 ID 对账——updateFromPoll用服务端队列的 ID 集合过滤掉已落地的 pending 项。这条 ID 约定贯穿 add/send/remove 三个端点是队列 UI 与后端保持一致的关键。5. 自动发送路径process_chain_end 扩展除了手动触发Agent Zero 还为队列提供了一条自动消费路径。扩展 extensions/python/process_chain_end/_50_process_queue.py 挂在“处理链结束”阶段ProcessQueue.execute第 11–39 行if self.agent.number ! 0: # 只对 agent0主 Agent生效 return context self.agent.context if mq.has_queue(context): asyncio.create_task(self._delayed_send(context))_delayed_send先以 0.1 秒步长轮询context.is_running()最长等待 60 秒防止挂死确认当前任务结束后调用mq.send_next(context)helpers/message_queue.pypop_firstsend_message的组合成功后以reasonmessage_queue_auto_send调用mark_dirty_for_context。从源码结构看这形成了清晰的优先级自动逐条消费是常态/message_queue_send是用户想“立即插队”或“一次性清空队列”时的主动干预口——两者共用同一套pop_*/send_message原语行为完全一致。6. 副作用与状态同步mark_dirty_for_context三条发送路径message_queue_send、message_queue_send_all以及自动路径的message_queue_auto_send在成功后都会调用mark_dirty_for_context。该函数定义在 helpers/state_monitor_integration.py语义是“标记该 context 的状态为脏”def mark_dirty_for_context(context_id: str, *, reason: str | None None) - None: from helpers.state_monitor import get_state_monitor get_state_monitor().mark_dirty_for_context(context_id, reasonreason)它把队列变化纳入helpers/state_monitor.py的状态监视体系驱动多端/多标签页的状态同步广播reason参数如message_queue_send_all则作为诊断信息记录“是什么操作弄脏了状态”。值得注意的是聚合分支的写法send_all_aggregated返回 0队列本来就是空的、被并发抢先消费时不会打脏标记避免无谓的同步风暴——这是一种轻量的幂等优化。7. DOX 工作法与验证建议message_queue_send.py.dox.md本身就是仓库工程约定的一部分值得提炼一文件一档api/目录刻意保持扁平每个.py配一个同名.dox.md分别持有“运行实现”与“持久化笔记”职责、契约、副作用、验证修改实现时必须同步改 DOX契约变更触发更新请求载荷、认证/CSRF 要求、响应形状、路由副作用、WebSocket 事件契约任一变化都要刷新 DOX验证策略DOX 的 Verification 一节写明“对变更行为运行端点专属或 API/WebSocket 测试若无聚焦测试对浏览器调用方做冒烟验证”。从仓库测试布局看tests/目录下未见以message_queue命名的专项用例该端点当前主要依赖行为最接近的 API 测试与前端冒烟路径来覆盖这与其“33 行、无内部状态”的低复杂度相匹配。8. 小结/message_queue_send是 Agent Zero 消息队列体系中负责“投递”的端点它以context定位会话用item_id/send_all两个开关区分“发指定条”“发队首”“聚合发全部”三种意图通过helpers/message_queue.py的 pop/communicate 原语完成投递并以mark_dirty_for_context触发前端状态同步配套的_sync_output截断视图与process_chain_end自动消费扩展共同构成了“忙时排队、闲时自动续发、随时手动插队”的完整体验。理解了这条链路上的 api/message_queue_add.py、api/message_queue_remove.py、helpers/message_queue.py 与 extensions/python/process_chain_end/_50_process_queue.py也就掌握了 Agent Zero 会话消息流控制的完整机制。【免费下载链接】agent-zeroAgent Zero AI framework项目地址: https://gitcode.com/GitHub_Trending/ag/agent-zero创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表