数据截至 (上游 commit 771b20c93db8)
内核桥:Jupyter 消息协议与流式执行(全项目最深)
30 秒导读: 沙箱里真正跑代码的是一个 Jupyter 内核。这一章讲的是夹在「FastAPI 路由」和「内核」之间的 那座桥——
messaging.py。它把用户的一次execute(code)翻译成 Jupyter 规定的execute_request消息发进内核,再把内核异步、乱序、分很多条吐回来的消息,重组成一条有序的输出流;还要在客户端悄悄断线时打断内核、在连接掉线时给在途执行补个「结束」信号防止永久卡死。全项目工程含量最高的一章。
这一章接在 02-server-routing.md(FastAPI 路由把请求交给哪个 WebSocket)之后。路由层做的最后一件事,就是把 ws.execute(...) 这个异步生成器塞进流式响应里(template/server/main.py:121-127)。本章讲的就是这个生成器内部的一切。
1. 这是什么(零基础也能懂)
先建立一个心智模型
想象你有一个「永远开着的 Python 终端」(其实是 Jupyter 内核)。你不能直接摸它,只能隔着一根网线跟它喊话。这根网线就是一条 WebSocket 长连接,喊话的「语言」是一套叫 Jupyter Wire Protocol(Jupyter 线协议) 的规矩:你发一条格式固定的 JSON 说「请执行这段代码」,内核就会陆陆续续回你好几条 JSON——"我开始忙了"、"这是一行 stdout"、"这是一张图"、"这是最终结果"、"我空闲了"。
messaging.py 干的活,就是把这套一来一回的喊话,包装成一个干净的 Python 异步生成器:
# 示意,非源码:调用方眼中的样子
async for chunk in ws.execute(code, env_vars, token):
# chunk 已经是结构化的 Stdout / Result / Error / EndOfExecution
send_to_client(chunk)
调用方完全不用知道 Jupyter 协议、不用管消息乱序、不用管连接掉线——桥全帮你抹平了。
它要解决的三个真问题
| 真问题 | 桥怎么接 |
|---|---|
| 内核回消息是异步 + 乱序 + 一次执行分很多条 | 后台任务 _receive_message 收所有消息,按 parent_msg_id 分发到对应执行的队列;_wait_for_result 从队列有序取出 |
| 客户端可能悄悄断线(网络挂了,没说再见) | _wait_for_result 每 5 秒发一个 keepalive 逼服务器写 socket,写失败才发现断了,进而打断内核(否则下一次执行被卡住) |
| WebSocket 掉线/重连时,在途执行会永远收不到结束 | 收消息循环退出时,给所有在途执行补一个 UnexpectedEndOfExecution,防止生成器永久 hang |
一句话直觉: 把它想成一个快递分拣中心。内核是四面八方发货的仓库(消息乱序到达),桥按每个包裹上的「订单号」(parent_msg_id)分拣进对应的传送带(队列),客户那头则从传送带上按顺序一件件取货,直到看见「本单发完」的标签(EndOfExecution)。
2. 顶层全景(它大概怎么转)
一次执行的生命周期
先给一句「怎么读这张图」:左边是一次 execute() 调用的主流程(串行发送),右边是永远在后台跑的收消息循环,两者靠中间那条 per-execution 队列解耦。
execute(code) ← 来自 FastAPI 路由(main.py:122)
│
│ async with _lock ← 串行化:同一时刻只放一次执行进内核
▼
┌───────────────┐ 注入 env vars(按首行缩进对齐)
│ 组装 complete │ 生成 msg_id + Execution(带自己的 queue)
│ _code │
└──────┬────────┘
│ _get_execute_request → JSON
│ ws.send(...) ← 失败则 reconnect 重发,最多 3 次
▼
┌──────────────────────┐ ┌───────────────────────────┐
│ _wait_for_result │ │ _receive_message (后台) │
│ while True: │◀──队列──│ async for msg in ws: │
│ queue.get(timeout) │ put() │ _process_message(msg) │
│ 超时→yield keepalive│ │ 按 msg_type 映射成 │
│ 拿到→yield 结构化 │ │ Stdout/Result/Error… │
│ 见 EndOfExec→break │ │ 按 parent_msg_id 分发 │
└──────────┬───────────┘ └────────────▲──────────────┘
│ │
▼ 内核 ws://localhost:8888
yield 给路由 → 客户端 /api/kernels/{id}/channels
│
(客户端断开 → CancelledError → shield(interrupt()) 打断内核)
│
▼
env_vars 清理:另起一个后台执行请求 _cleanup_env_vars
部件一句话职责
| 部件(符号) | 干什么 | 位置 |
|---|---|---|
ContextWebSocket | 一个 context 一条长连接的门面,握着 ws、收消息任务、执行表 | messaging.py:60 |
Execution | 一次执行的信箱:一条 asyncio.Queue + errored/input_accepted 标志 | messaging.py:42 |
connect / reconnect | 建连并起后台收消息任务 | messaging.py:75-102 |
execute | 主流程:串行化、注入 env、发送带重连、流式 yield、断开打断 | messaging.py:318 |
_wait_for_result | 从某次执行的队列有序取输出,带 keepalive 超时 | messaging.py:254 |
_receive_message | 后台读所有内核消息,掉线时给在途执行补收尾 | messaging.py:422 |
_process_message | 把每种 Jupyter msg_type 映射成结构化对象入队 | messaging.py:445 |
_get_execute_request | 按 Jupyter 协议构造 execute_request JSON | messaging.py:120 |
interrupt | 走 Jupyter REST 打断内核 | messaging.py:104 |
change_current_directory | 各语言用不同 magic 切工作目录 | messaging.py:289 |
关键设计:每个 context 一条长连接。 __init__ 只是把 url 拼成 ws://localhost:8888/api/kernels/{context_id}/channels(messaging.py:70),真正建连在 connect()——一次建连后 _receive_task 就一直在后台跑,后续所有执行复用同一条连接(messaging.py:90-102)。