WebSocket 实时更新
第3周收官。执行过程如何实时推送到前端画布——你看到节点一个个亮起、数据流动的原理。读 api/ws_api.py + conn_manager.py。
👶 小白:执行是在后台进程里跑的,前端浏览器怎么就能"看到"节点一个个亮起?
👨🏫 老师:靠一条转播链。执行进程每完成一步就发一条事件到 Redis(后台演播),WS 进程订阅到后推给正在看这场"直播"的你的浏览器,画布收到就点亮对应节点。执行进程和你的浏览器不直接通话——中间隔着 Redis 和 WS 进程转播,这样执行只管埋头干、推送由专门的直播间负责,互不拖累。
实时进度需求
你在画布上点"运行"后,想实时看到:哪个节点在跑、哪个完成了、数据是什么、有没有报错——而不是干等到全部跑完才看结果。这需要后端把执行进度实时推给前端。HTTP 请求-响应模型做不到"服务器主动推",所以用 WebSocket。
独立 WS 进程
WebSocket 是独立进程/独立 app(api/ws_api.py:55),不跟 REST 混在一起。唯一端点 @app.websocket("/ws")(:230)。(Day 04 提过 REST 服务显式 ws="none",把 WS 交给这个独立进程。)
鉴权
authenticate_websocket()(ws_api.py:75)从 query 参数取 JWT token,解析出 user_id,失败就 websocket.close(code=4001/4002/4003)。
ws://.../ws?token=xxx)。连上后立刻验 token、解析出是哪个用户——验不过直接关连接(4001/4002/4003 是自定义关闭码,区分不同失败原因)。鉴权是必须的——否则任何人都能订阅别人的执行进度(泄漏)。订阅消息路由
客户端连上后发消息,_MSG_HANDLERS(ws_api.py:219)字典分发四种方法:
HEARTBEAT:心跳 pong(保活)。SUBSCRIBE_GRAPH_EXEC:订阅单次执行的进度。SUBSCRIBE_GRAPH_EXECS:订阅某个图的所有执行。UNSUBSCRIBE:取消订阅。
graph_exec_id = abc123 后,前端通过 WS 发:{"method": "subscribe_graph_execution", "data": {"graph_exec_id": "abc123"}}→
_MSG_HANDLERS(ws_api.py:219)按 method 分发到 handle_subscribe(ws_api.py:109)→ 调 connection_manager.subscribe_graph_exec(ws_api.py:142)登记订阅。此后
abc123 的每条进度事件都会自动推到这个连接,直到你发 unsubscribe。Redis → WebSocket 转发链
ConnectionManager(conn_manager.py:219)管理所有 WS 连接:用户连上走 connect_socket(:229);订阅某次执行走 subscribe_graph_exec(:249)——它去 Redis Pub/Sub 订阅那条频道;后台 _pump() 一直 listen(),收到 Redis 事件就 _forward_exec_event(:315)挑出订阅了它的 WebSocket 发出去。executor 只管"发到 Redis 频道",谁听、谁转发它不管——彻底解耦。推送机制在 api/conn_manager.py 的 ConnectionManager。订阅时 subscribe_graph_exec(:249)打开一个 Redis Pub/Sub 订阅;_pump()(:115)listen() 收到 Redis 事件后 _forward_exec_event(:315)转发给对应 WebSocket 客户端。
完整推送链路(一图流)
publish,到你眼前的动画——中间隔着 Redis 中转和 WS 转发,但对你透明。历史数据(刷新页面重进)则走 REST 从 DB 拉(Day 18),实时 + 历史拼成完整体验。为什么用 WebSocket 而非 SSE
OpenHands 教程(Day 16)里也讨论过这个。AutoGPT 用 WebSocket 而非 SSE,原因:
- 双向:前端不只"收进度",还要"发订阅/取消订阅/心跳"——需要双向通信。SSE 只能服务器单向推。
- 多路复用:一个 WS 连接能订阅多个执行、动态增减订阅——比每个执行开一个 SSE 连接高效。
🎓 第 3 周收官 + 动手
第 3 周(Day 11-15)你已吃透执行引擎
- Day 11 Graph 模型:三层、计算字段、校验、版本/子图
- Day 12 执行引擎:MQ+线程池+事件总线、调度循环、分布式锁
- Day 13 数据流:输出喂下游、攒齐输入才入队、fan-in
- Day 14 Scheduler:定时、持久化、cron 兼容、自愈
- Day 15 WebSocket:实时进度推送链路
你已理解"图怎么被执行、进度怎么实时展示"。下周(Day 16-20)进入 平台底座 + 起源:计费、凭证/OAuth、REST API、经典 AutoGPT、收官。
✋ 动手
P=autogpt_platform/backend/backend
grep -n 'def authenticate_websocket\|_MSG_HANDLERS\|websocket(' $P/api/ws_api.py | head
grep -n 'def subscribe_graph_exec\|def _pump\|_forward_exec_event' $P/api/conn_manager.py