数据截至 (上游 commit a4eaba4a56f9)
核心机制 — 子进程传输层与双向控制协议
本章是整个 SDK 的心脏。讲两件事:(1)
SubprocessCLITransport怎么把配置变成一个跑起来的claude子进程、怎么读写管道;(2)Query怎么在这条管道上跑一套双向 RPC 协议,让 CLI 能"回头问"你的 Python 代码。
1. 大图:两类消息共享一条管道
SDK 和子进程之间只有 stdin/stdout 两根水管,但上面跑着两类语义完全不同的 JSON 行:
| 类别 | 方向 | 干什么 | 例子 type 值 |
|---|---|---|---|
| SDK 消息 | CLI → SDK | 对话数据流,产出给你 | assistant user system result |
| 控制消息 | 双向 | RPC:一方请求、另一方回响应 | control_request control_response control_cancel_request |
Query._read_messages 就是那个"分拣员":读到一行先看 type,是控制消息就内部处理掉(不产出给你),是 SDK 消息才送进产出队列。
怎么读下面这张图:实线是数据流向,①②③ 是一次"CLI 反问权限"的时序。
你的 Python Query 读循环 claude 子进程
│ │ │
│ write(user JSON) ───────────┼────────── stdin ────────▶│
│ │ │
│ │◀──── stdout: assistant ─────│ 数据流
│◀── yield AssistantMessage ───│ │
│ │ │
│ │◀─ stdout: control_request ①─│ "Bash 能用吗?"
│◀─ can_use_tool 回调 ②────────│ │
│── 返回 allow/deny ───────────▶│ │
│ │─ stdin: control_response ③ ▶│
│ │ │
│ │◀──── stdout: result ────────│ 这轮结束
│◀──── yield ResultMessage ────│ │
2. 传输层:SubprocessCLITransport
2.1 它要解决的小问题
"把 ClaudeAgentOptions 这一大坨配置,变成一个真正跑起来的 claude 进程,并可靠地读写它的管道、退出时不留僵尸进程。"
2.2 找 CLI:三级 fallback
启动前得先找到 claude 二进制。查找顺序(subprocess_cli.py:247 _find_cli):
① 捆绑的 CLI _bundled/claude(.exe) ← 随 pip 包一起装,默认走这条
② PATH 里的 shutil.which("claude")
③ 常见安装位 ~/.npm-global/bin、/usr/local/bin、~/.local/bin …
找不到 → 抛 CLINotFoundError,附带安装指引
捆绑优先是这个包的特点:README 说"CLI 随包捆绑,无需另装"。你也能用 ClaudeAgentOptions(cli_path=...) 指定路径。
2.3 拼命令行:配置 → CLI flags
_build_command(subprocess_cli.py:562)是一长串"如果设了这个 option 就 append 那个 flag"。骨架永远是:
claude --output-format stream-json --verbose ... --input-format stream-json
几个值得注意的翻译(都在 _build_command 里):
system_prompt分三种:纯字符串 →--system-prompt;{type:file}→--system-prompt-file;{type:preset, append}→--append-system-prompt(subprocess_cli.py:568-579)。mcp_servers里的 SDK 类型服务器要剥掉instance字段再传--mcp-config(subprocess_cli.py:657-679)——因为那个 Python 对象没法序列化过管道,实例留在本进程,靠控制协议回调。skills会顺带补allowed_tools(注入Skill或Skill(name))和默认setting_sources=["user","project"](_apply_skills_defaults,subprocess_cli.py:519)——省得你手动配两处。- agents 不走命令行,而是通过 initialize 请求发(
subprocess_cli.py:716注释明说)。
2.4 开进程 + 环境变量处理
connect()(subprocess_cli.py:787)用 anyio.open_process 开子进程。环境变量合并有两个巧妙处(subprocess_cli.py:809-814):
- 过滤掉
CLAUDECODE:防止 SDK 拉起的子进程误以为自己跑在一个 Claude Code 父进程里(issue #573)。 - 强制设
CLAUDE_CODE_ENTRYPOINT=sdk-py和CLAUDE_AGENT_SDK_VERSION,但你的options.env能覆盖前者。
还有一层可选的 OTEL 追踪上下文注入(subprocess_cli.py:820-841):如果装了 opentelemetry-api 且有活跃 span,就把 W3C traceparent 注进子进程环境,让 CLI 的 span 挂在你的分布式追踪下。失败绝不影响 connect(best-effort)。
2.5 防僵尸进程
模块级维护一个 _ACTIVE_CHILDREN 集合,atexit 注册 _kill_active_children(subprocess_cli.py:50-60):Python 进程退出时给所有活着的子进程发 SIGTERM。这对应 TypeScript SDK 的父退出清理,防止调用方崩溃时漏掉一堆 claude 进程。
2.6 读 stdout:为什么要"投机式"缓冲
_read_messages_impl(subprocess_cli.py:1081)不是简单地"一行一 JSON"。因为 TextReceiveStream 会把长行截断,所以它:
累积 json_buffer ← 每来一段就拼上
尝试 json.loads(json_buffer)
成功 → yield,清空 buffer
失败(JSONDecodeError) → 继续累积,不报错
buffer 超过 max_buffer_size → 抛 CLIJSONDecodeError
两个防御细节:
- 完整行不以
{开头(如[SandboxDebug]调试输出)解析时直接跳过,免得污染消息流(subprocess_cli.py:205-208,issue #347)。 - 子进程非零退出码 → 造一个
ProcessError抛出(subprocess_cli.py:1136-1142)。
2.7 写 stdin 与优雅关闭
写用一把 anyio.Lock(_write_lock)串行化,并在锁内做所有"能不能写"的检查,避免和 close()/end_input() 的 TOCTOU 竞态(subprocess_cli.py:1043-1067)。
关闭时先给 5 秒优雅期再考虑 SIGTERM(subprocess_cli.py:1004-1026):子进程收到 stdin EOF 后需要时间刷会话文件,直接 SIGTERM 会打断写、丢最后一条 assistant 消息(issue #625)。超时才 terminate(),再超时才 kill()。
3. 控制协议:Query 类
3.1 它要解决的小问题
光有单向数据流不够。有些事 CLI 必须"回头问"SDK:这个工具准不准用?跑一下你那个 Python 工具、给我结果?这就需要一套在同一管道上跑的双向 RPC。Query 就是这套协议的路由中枢(_internal/query.py:102)。
3.2 请求/响应怎么配对
核心是"每个请求一个唯一 id + 一个 anyio.Event"的经典异步 RPC 模式:
发起方 _send_control_request(request):
request_id = f"req_{counter}_{随机hex}"
pending_control_responses[request_id] = Event() ← 挂一个等待点
写 {"type":"control_request", request_id, request} 到 stdin
with fail_after(timeout): await event.wait() ← 阻塞直到读循环唤醒
取 pending_control_results[request_id] 返回
读循环收到 control_response:
按 response.request_id 找到那个 Event
存结果到 pending_control_results
event.set() ← 唤醒发起方
依据:_send_control_request(query.py:598)、_read_messages 里 control_response 分支(query.py:318-330)。超时会清理挂起状态并抛异常(query.py:640-643)。
3.3 读循环:分拣五类消息
_read_messages(query.py:308)是 Query 的主循环,按 type 分拣:
| type | 处理 | 出处 |
|---|---|---|
control_response | 唤醒对应挂起请求 | query.py:318 |
control_request | 反向请求:spawn 一个 handler 去处理 | query.py:332 |
control_cancel_request | 取消一个在途的 handler | query.py:340 |
transcript_mirror | 会话镜像帧:交给 batcher,不产出 | query.py:348 |
其它(result/assistant/…) | 送进产出队列 | query.py:399 |
3.4 反向请求:CLI 让 SDK 干活
_handle_control_request(query.py:469)处理 CLI 发来的三种 subtype:
| subtype | 干什么 | 调用你的什么 |
|---|---|---|
can_use_tool | 问某工具能不能用 | can_use_tool 回调(见第 03 章) |
hook_callback | 触发某个钩子 | 你注册的 hook 函数 |
mcp_message | 转发一条 MCP JSONRPC 给进程内服务器 | 你的 @tool 函数 |
处理完把结果包成 control_response(subtype success 或 error)写回 stdin。注意:如果处理途中收到 control_cancel_request 导致取消,不写响应——CLI 已经放弃这个请求了(query.py:582-585)。
3.5 进程内 MCP:桥接真实会话
_handle_sdk_mcp_request(query.py:645-674)如今是个薄路由:每个 SDK MCP 服务器在 Query.__init__ 里就包成了一个 SdkMcpBridge(query.py:153-157),消息到来时直接交给 bridge.handle()。真正的分派在 SdkMcpBridge(sdk_mcp_bridge.py)——用 MCP 官方的内存传输对跑一个真实的 Server.run 会话,initialize / tools / resources / prompts / ping 全部由 mcp 库自身分派(模块 docstring sdk_mcp_bridge.py:1-23 道明设计)。旧版"手写 JSONRPC 分派、等上游补 Transport 再重构"的 TODO 已还清。
3.6 initialize 握手
initialize(query.py:231)是连接后发的第一条控制请求,把这些"没法走命令行"的配置递给 CLI:
hooks:为每个钩子回调分配hook_{id},建立 id → 函数映射(query.py:245-259)——之后 CLI 用 id 回调,SDK 靠这张表找函数。agents、excludeDynamicSections、skills(仅当是显式 list 时才发,"all"和省略等价,query.py:272)。
超时用 initialize_timeout(默认至少 60 秒),因为 MCP 服务器启动可能慢。
3.7 一个精妙处:把"exit code 1"换成真错误
当 CLI 报一个 result 且 is_error=True(如 error_max_turns),它会故意非零退出(为了 shell 脚本消费者)。这样传输层就会抛一个只带"exit code 1"、毫无信息量的 ProcessError。Query 记住最后那条 error result 的完整消息(_last_error_result,query.py:384-387),在抛出前把空洞的 ProcessError 替换成携带原始 result 与退出码的类型化 ResultError(query.py:414-428;ResultError 定义在 _errors.py:56)。这直接对齐 TypeScript SDK 的 Query.ts 逻辑。
4. 巧妙之处小结
- 一条管道跑两套语义:数据流和 RPC 复用同一 stdin/stdout,靠
type字段分拣。 - 永远 streaming:即使
query(str)也内部走流式,统一了两个入口的实现,让 agents/hooks 总能通过 initialize 发出去。 - 投机式 JSON 缓冲:容忍被截断的长行,而不是假设"一行一 JSON"。
- 优雅关闭 5 秒窗口:换来最后一条消息不丢。
5. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 传输实现 | src/claude_agent_sdk/_internal/transport/subprocess_cli.py | SubprocessCLITransport |
| 找 CLI | 同上 | _find_cli _find_bundled_cli |
| 拼命令行 | 同上 | _build_command _apply_skills_defaults |
| 读 stdout | 同上 | _read_messages_impl |
| 写与关闭 | 同上 | write close end_input |
| 防僵尸 | 同上 | _kill_active_children _ACTIVE_CHILDREN |
| 协议中枢 | src/claude_agent_sdk/_internal/query.py | Query |
| 读循环分拣 | 同上 | _read_messages |
| 发控制请求 | 同上 | _send_control_request |
| 处理反向请求 | 同上 | _handle_control_request _handle_sdk_mcp_request |
| 握手 | 同上 | initialize |