拓冰建站拓冰建站
首页 / 资讯中心 / 正文

agno AgentOS 长时运行工作流韧性实战:WebSocket 断线重连、事件 Replay 与 Catch-Up 追赶机制

agno AgentOS 长时运行工作流韧性实战WebSocket 断线重连、事件 Replay 与 Catch-Up 追赶机制【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno工作流Workflow以分钟乃至小时为单位在后台运行时WebSocket 客户端会随时断开如何在重连后不丢失任何中间事件、如何把断线期间产生的错过事件一次补齐、如何从已结束的运行中完整回放全部历史是本篇要解决的核心问题。本文以 cookbook/04_workflows/06_advanced_concepts/long_running 目录下三个可直接运行的可执行示例为主线结合 agno 事件流event stream底层实现讲解基于event_index的断点续传协议读完即可照着把用户刷新页面 / 掉线重连场景做成工作流前端接入的标准方案。目录定位一个专注于长时运行韧性的可运行示例集该目录属于工作流高级概念cookbook/04_workflows/06_advanced_concepts之下README 明确其定位为Scopecookbook/04_workflows/06_advanced_concepts/long_running下**可运行Runnable**的工作流示例集合Files共三个示例分别覆盖三类典型的长时运行故障恢复场景文件演示主题核心场景disruption_catchup.pyDisruption Catchup中断追赶重连运行中的工作流携带last_event_indexNone获取全部历史事件events_replay.pyEvents Replay事件回放重连已完成的工作流验证last_event_index被忽略、事件从头完整重放websocket_reconnect.pyWebSocket Reconnect重连模拟用户离开页面 3 秒再重连验证订阅、断线、补发错过事件全链路从源码结构看三者都是纯 WebSocket 客户端测试程序通过websockets库连到本地 AgentOS 服务端的ws://localhost:7777/workflows/ws向名为content-creation-workflow的内容创作工作流下发start-workflow/reconnect指令再按事件类型与索引号做端到端正确性校验。事件协议里的几个核心字段三个脚本反复解析并校验以下字段这是理解整个机制的关键event事件类型如WorkflowStarted、WorkflowCompleted、WorkflowError以及三条控制类通知catch_up、replay、subscribedevent_index客户端可见的、严格单调递增的事件序号从 0 开始是断点续传的游标run_id一次工作流运行的全局唯一 ID重连时用它定位要订阅的运行missed_eventscatch_up 通知携带断线期间错过的事件数量current_event_count当前已产生的事件总数用于进度对齐status/total_eventsreplay 通知携带回放的状态与总事件数。运行前置条件README 给出的环境要求与文件头注释一致共三条激活 demo 虚拟环境使用.venvs/demo/bin/python解释器运行示例TEST_LOG.md 亦记录所有脚本均以.venvs/demo/bin/python执行加载 API 密钥通过direnv allow加载需要本地.envrc文件启动本地 AgentOS 服务部分示例依赖运行在http://localhost:7777的 AgentOS 工作流服务端即ws://localhost:7777/workflows/ws的宿主。脚本文件头提示先启动对应的工作流服务再运行测试脚本额外依赖websockets库若缺失脚本会提示uv pip install websockets后退出。示例一WebSocket 断线重连与错过事件追赶websocket_reconnect.pywebsocket_reconnect.py 封装了一个WorkflowWebSocketTester类把用户在长时运行期间离开页面再回来的真实体验抽象为两个阶段。Phase 1启动工作流并接收初始事件连接ws://localhost:7777/workflows/ws后先收到服务端第一条欢迎消息随后发送启动指令await websocket.send( json.dumps( { action: start-workflow, workflow_id: content-creation-workflow, message: Research and create content plan for AI agents, session_id: test-session-123, } ) )客户端持续接收事件实时记录两个关键状态并维护last_event_index作为后续续传游标if run_id in data and not self.run_id: self.run_id data[run_id] if event_index in data: self.last_event_index data[event_index]当累计收到 20 个事件max_initial_events 20时脚本主动模拟断线——记录下最后一次拿到的事件序号后退出async with块WebSocket 随之关闭若在此之前工作流已结束WorkflowCompleted/WorkflowError同样终止接收。Phase 2重连并补齐错过事件脚本用asyncio.sleep(3)模拟用户离开页面 3 秒期间工作流服务端在后台继续产生事件。随后重新建立连接并发送reconnect指令把 Phase 1 保存的run_id与last_event_index一并回传await websocket.send( json.dumps( { action: reconnect, run_id: self.run_id, last_event_index: self.last_event_index, workflow_id: content-creation-workflow, session_id: test-session-123, } ) )重连后的事件流包含三类关键通知脚本分别处理catch_up通知携带missed_events错过的事件数与current_event_count随后服务端把断线期间积压的事件一次性补发replay通知当运行已结束、改走完整回放路径时出现subscribed通知追赶完成后客户端正式转为订阅新事件状态status与current_event_count用于确认对齐点。脚本对后续每个事件用MISSED/NEW标记区分补发与新到直到WorkflowCompleted/WorkflowError结束。最终_print_summary()会对收到事件做三种验证事件类型分布统计、首尾 event_index 跨度以及最严格的索引连续性检查——用set(range(min, max 1))与真实索引集合求差有缺口即报Gaps in event_index。示例二中断完全追赶——全量历史事件补发disruption_catchup.pydisruption_catchup.py 验证的是追赶机制的上限场景重连一个仍在运行的工作流时显式传last_event_indexNone语义等价于我什么事件都不要漏从 0 号事件开始全部给我。Phase 1 等待窗口与示例一同理先启动工作流并接收事件但这次只收3 个事件就主动断开制造一个丢失了大量事件的窗口if event_count max_events: print(f\nDisconnecting after {event_count} events...) break随后asyncio.sleep(3)等待工作流继续产出事件。Phase 3last_event_indexNone的全量追赶重连后发送reconnect指令last_event_index显式置空await websocket.send( json.dumps( { action: reconnect, run_id: run_id, last_event_index: None, workflow_id: content-creation-workflow, session_id: full-catchup-test, } ) )服务端行为差异在此体现收到last_event_indexNone后不再做增量对齐而是从事件 0 开始回放。客户端侧判定逻辑把到达事件分为两组——收到subscribed通知之前的事件归入catchup_events之后的归入new_eventsif event_type catch_up: print(f missed_events: {data.get(missed_events)}) print(f current_event_count: {data.get(current_event_count)}) if event_index is not None: if not got_subscribed: catchup_events.append(data) # 追赶段 else: new_events.append(data) # 订阅后的新事件脚本最终断言三项关键事实收到catch_up通知got_catch_up为真追赶段首个事件event_index为 0first_index 0否则报should be 0证明拿到的是完整历史而非增量追赶段事件索引无缺口且在subscribed之后仍有新事件持续到达证明先补历史、后接实时的衔接没有缝隙。脚本结尾的 Key Takeaway 一句话点题Send last_event_indexNone to get ALL events from start, even when reconnecting to a RUNNING workflow——对运行中工作流做全量追赶就是要传None。示例三已完成运行的整段事件回放events_replay.pyevents_replay.py 针对的是与前两个示例相反的状态工作流已经结束WorkflowCompleted此时任何增量续传语义都失去意义重连应当触发全量事件回放。Phase 1等待工作流自然完成启动工作流后持续接收直到WorkflowCompleted期间用total_events max(total_events, data[event_index] 1)记录完整运行的总事件数作为后续回放正确性的基准。Phase 2携带无效游标重连验证游标被忽略重连时故意发送一个过期的游标last_event_index10用于验证服务端对已完成运行的处理策略await websocket.send( json.dumps( { action: reconnect, run_id: run_id, last_event_index: 10, # 应当被忽略 workflow_id: content-creation-workflow, session_id: replay-test-session, } ) )随后收到replay通知内含status、total_events与说明性message。校验逻辑如下回放首个事件event_index必须为 0回放事件总数必须等于 Phase 1 记录的总数len(replay_events) total_events从而证明last_event_index对已完成运行被忽略回放序列无索引缺口。该脚本与前两例组合恰好覆盖了断点续传协议的三个分支运行中 无游标 → 全量追赶、运行中 有游标 → 增量追赶 订阅、已结束 任意游标 → 全量回放。底层原理event_index单调非连续契约与 tail/replay 竞态三个示例所依赖的语义并非黑盒魔法其契约在 agno 事件流基类中写得很清楚。libs/agno/agno/os/event_streams/base.py 开头的实现注释给出了四条对客户端至关重要的约定event_index是客户端可见的 per-run 索引strictly increasing, but NOT guaranteed gapless——严格递增但不保证无缺口。生产者producer在分配索引与发布事件之间崩溃可能留下永久缺口。因此客户端必须以last_event_index续传且必须容忍缺口运行终止由run status判定绝不能依赖收到某个特定索引号来判断结束。这正是三个示例都做缺口检查并打印、但以WorkflowCompleted为真正退出条件的原因tail()必须自己处理订阅/回放竞态不能漏掉调用方完成 replay 之后、真正开始 tail 之前到达的事件也不要求调用方用锁协调。对应到客户端视角就是catch_up/replay补发段与subscribed实时段必须无缝衔接——示例二据此判定衔接正确性tail()必须在 run 到达终态时终止即使生产者已死、没来得及写终态标记实现也必须在空闲时复查 run status而不是无限阻塞等待。从 router.py 附近的代码看服务端对reconnect指令的解析同样基于这一游标模型last_event_index被注释明确为客户端已收到的最后一个事件的 0-based 索引重连处理流程会区分运行仍在缓冲区内回放last_event_index之后的增量事件后转入 tail与运行已完成直接全量回放两条路径与示例二、示例三观测到的客户端行为一一对应。此外 router 中还存在对last_event_indexNone场景的兜底replay 起始索引取None时即从全部事件开始解释了None触发全量补发的机制。对于多副本/分布式部署事件流是可插拔的默认内存实现在单进程内缓冲Redis Streams 等分布式实现让事件在任何容器都可读客户端可以在执行运行的副本之外恢复流。若需追溯事件流替换机制可看 event_streams/base.py 的BaseEventStream抽象类与其通过set_event_stream()替换全局实例的注册模式。运行与验收从脚本到测试日志三个脚本均可独立运行入口统一为asyncio.run(main())且都做了连接失败友好处理捕获ConnectionRefusedError并提示AgentOS server 是否在运行、如何启动。开始执行前会打印前置条件并等待 2 秒方便用户核对环境。TEST_LOG.md 记录了三个脚本在同一 demo 环境下的验收结果三者均以.venvs/demo/bin/python执行状态均为 PASS。需要注意其运行模式为startup validation only启动验证模式timeout 2 秒即 CI 仅验证脚本能正常导入、启动并打印 Starting test in 2 seconds... 后即被终止——完整的断线追赶行为需要在本地 AgentOS 服务可用时人工运行验证这一点与示例文件头的前置说明一致。小结与工程实践要点综合三个示例与底层契约接入长时运行工作流事件流时应遵循以下要点永远以event_index为游标做重连续传并记住严格递增、不保证无缺口——客户端逻辑必须以 run 终态WorkflowCompleted/WorkflowError判定结束不能依赖索引连续三种重连语义按运行状态区分运行中 已知游标走增量追赶catch_up→subscribed运行中 last_event_indexNone走全量追赶已结束的运行无论游标如何都走全量回放replay此时游标被忽略以run_id定位运行、以session_id区分会话断线重连只关心run_id是否还活着做好本地可观测验证三个示例内置的缺口检查 事件类型统计 首尾索引核对三段式校验可以直接复用到你自己的前端重连逻辑测试中理解缓冲/回放的边界事件缓冲有保留期与容量长时运行若远超缓冲上限追赶上界取决于服务端事件流实现内存 or Redis 等生产环境应按需选择分布式事件流并配合持久化兜底。如果你正基于 agno 开发带实时进度页面的长任务内容创作、批量处理、多智能体编排这三个示例是目前最直接的断点续传 事件回放参考实现。【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门