Day 15 / 共 20 天 · 第 3 周 执行引擎(收官)

WebSocket 实时更新

第3周收官。执行过程如何实时推送到前端画布——你看到节点一个个亮起、数据流动的原理。读 api/ws_api.py + conn_manager.py

📍 你在整门课的位置 · 第 3 周 执行引擎(收官)
D14 Scheduler D15 WebSocket D16 Credits 计费 D17 OAuth
💡 今天的类比世界观:实时推送 = 一场"直播" 你在画布上看到节点一个个亮起,靠的是一条直播链路:独立 WS 进程 = 专门的直播间(不占用 API 的活);鉴权 = 进直播间先验票订阅路由 = 你只收自己关注的那个频道Redis → WebSocket 转发 = 后台把执行消息经 Redis 广播,再推到你的屏幕WS vs SSE = 双向电话 vs 单向收音机。今天都用"直播"来想。

👶 小白:执行是在后台进程里跑的,前端浏览器怎么就能"看到"节点一个个亮起?

👨‍🏫 老师:靠一条转播链。执行进程每完成一步就发一条事件到 Redis(后台演播),WS 进程订阅到后推给正在看这场"直播"的你的浏览器,画布收到就点亮对应节点。执行进程和你的浏览器不直接通话——中间隔着 Redis 和 WS 进程转播,这样执行只管埋头干、推送由专门的直播间负责,互不拖累。

L01

实时进度需求

你在画布上点"运行"后,想实时看到:哪个节点在跑、哪个完成了、数据是什么、有没有报错——而不是干等到全部跑完才看结果。这需要后端把执行进度实时推给前端。HTTP 请求-响应模型做不到"服务器主动推",所以用 WebSocket。

HTTP 的局限 普通 HTTP 是"客户端问、服务器答"——服务器不能主动给客户端发消息。但执行进度是服务器端不断产生的、需要主动推给你。WebSocket 是"全双工长连接"——建立后双方都能随时发消息。于是执行引擎每有进度,就通过 WebSocket 推给前端,画布实时更新。这是所有"实时看进度"UI 的技术基础。
L02

独立 WS 进程

WebSocket 是独立进程/独立 appapi/ws_api.py:55),不跟 REST 混在一起。唯一端点 @app.websocket("/ws"):230)。(Day 04 提过 REST 服务显式 ws="none",把 WS 交给这个独立进程。)

为什么 WS 要独立进程? WebSocket 是长连接——每个在线用户占一个连接、可能挂很久。REST 是短请求——来了就走。两者的资源特性、伸缩需求完全不同:WS 进程要扛"大量长连接",REST 进程要扛"高频短请求"。分开进程,各自按自己的特性优化和伸缩(比如 WS 进程可以配更多内存扛连接数)。不同性质的负载分进程——和 Day 04 拆执行/REST 一个道理。
L03

鉴权

authenticate_websocket()ws_api.py:75)从 query 参数取 JWT token,解析出 user_id,失败就 websocket.close(code=4001/4002/4003)

WebSocket 怎么鉴权?(和 HTTP 不同) HTTP 每个请求带 Authorization header。但 WebSocket 建连时浏览器 API 不方便带自定义 header,所以常见做法是把 token 放在连接 URL 的 query 参数里ws://.../ws?token=xxx)。连上后立刻验 token、解析出是哪个用户——验不过直接关连接(4001/4002/4003 是自定义关闭码,区分不同失败原因)。鉴权是必须的——否则任何人都能订阅别人的执行进度(泄漏)。
L04

订阅消息路由

客户端连上后发消息,_MSG_HANDLERSws_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_HANDLERSws_api.py:219)按 method 分发到 handle_subscribews_api.py:109)→ 调 connection_manager.subscribe_graph_execws_api.py:142)登记订阅。
此后 abc123 的每条进度事件都会自动推到这个连接,直到你发 unsubscribe
"订阅"是什么模式? 前端说"我要看执行 xxx 的进度"(SUBSCRIBE),服务器就把这个执行的所有事件推给它,直到它 UNSUBSCRIBE。这是发布-订阅:客户端订阅感兴趣的"话题"(某次执行/某个图),服务器只推它订阅的。心跳(HEARTBEAT)则是长连接的保活机制——定期 ping-pong,确认连接还活着(否则网络中断了双方都不知道)。
L05

Redis → WebSocket 转发链

🤔 痛点:跑图的进程和连着你浏览器的进程根本不是同一个,事件怎么送到你眼前? executor 进程(Day 12)在某台机器上跑图、产生进度事件;你的浏览器连的是另一台机器上的 WS 进程。executor 压根不知道"你"连在哪个 WS 进程上——它怎么把"节点 3 完成了"这条消息,精确送到正看着这张图的那个人的屏幕?
💡 本质:Redis 当"中间人广播",WS 进程只转发自己订阅的那些执行 ConnectionManagerconn_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.pyConnectionManager。订阅时 subscribe_graph_exec:249)打开一个 Redis Pub/Sub 订阅_pump():115listen() 收到 Redis 事件后 _forward_exec_event:315)转发给对应 WebSocket 客户端。

读法:关键链条:executor 进程 publish 事件到 Redis(Day 12 的事件总线)→ WS 进程订阅 Redis → 收到后转发给前端 WebSocket。Redis 是 executor 和 WS 两个进程之间的"中间人"。
为什么要经过 Redis 中转? 因为 executor(跑图的)和 WS 进程(连着前端的)是两个不同的进程,甚至不同机器。executor 产生的事件怎么送到"连着你浏览器的那个 WS 进程"?用 Redis Pub/Sub 做中转:executor 发到 Redis 频道,所有 WS 进程都订阅,谁连着相关用户谁就转发。这解耦了"产生事件"和"推送事件"——executor 不用知道用户连在哪个 WS 进程上,只管发到 Redis。Redis 事件总线是跨进程通信的桥。
两个进程 + Redis 中间人 executor 进程 跑图,产生进度 publish 事件 Redis Pub/Sub 频道 exec/{id} 广播给订阅者 WS 进程 ConnectionManager :219 _pump listen() :115 _forward_exec_event :315 前端画布:节点亮起 publish 订阅收到 WS 转发
executor 只 publish 到 Redis 频道;每个 WS 进程订阅并只转发自己连着的用户所订阅的执行——Redis 让"产生事件"和"推送事件"彻底解耦。
L06

完整推送链路(一图流)

executor节点状态变化 Redispublish 到频道 WS 进程订阅收到 WebSocket转发 前端画布节点亮起
这就是你在画布上看到"节点一个个变绿、数据在连线上流动"的完整技术链路。从 executor 的一行 publish,到你眼前的动画——中间隔着 Redis 中转和 WS 转发,但对你透明。历史数据(刷新页面重进)则走 REST 从 DB 拉(Day 18),实时 + 历史拼成完整体验。
L07

为什么用 WebSocket 而非 SSE

OpenHands 教程(Day 16)里也讨论过这个。AutoGPT 用 WebSocket 而非 SSE,原因:

  • 双向:前端不只"收进度",还要"发订阅/取消订阅/心跳"——需要双向通信。SSE 只能服务器单向推。
  • 多路复用:一个 WS 连接能订阅多个执行、动态增减订阅——比每个执行开一个 SSE 连接高效。
对比 OpenHands 的选择 有意思的是:OpenHands 前端也用 WebSocket(Day 16),理由一致(双向)。凡是"前端既要收实时数据、又要发控制指令"的场景,WebSocket 都是首选;纯服务器单向推(如股票行情)才用 SSE。三个 Agent 平台(AutoGPT/OpenHands 的实时通道都是 WS)殊途同归——因为交互式 Agent 天然需要双向实时通信。
L08

🎓 第 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
下周预告 · Day 16:进入平台底座——Credits 积分计费深入:两阶段扣费(预扣+对账)、原子 SQL 保证余额正确、Stripe 充值闭环、余额预警。这是计量 SaaS 的钱袋子。
← Day 14 Scheduler Day 16 · Credits →