05-IPC 与网络通信盲区梳理
IPC 与网络通信盲区梳理
概述
本篇梳理在学习 MiniClaude 项目过程中,围绕进程间通信(IPC)、TCP 网络通信、双进程架构、消息分流暴露的知识盲区。每条盲区包含:①原来的困惑/错误理解 ②正确解释 ③代码示例 ④延伸知识。
学习者在 Java 阶段习惯了"线程 + RPC 框架"的并发模型,进入 Python asyncio + 裸 TCP + JSON-RPC 的世界后,对"回调函数到底谁在调""为什么一条 TCP 上要跑两个协程""Future 怎么把请求和响应配对"等机制层问题反复卡壳。本篇把这些盲区按"基础回调 → 跨进程架构 → 消息分流 → 并发与生命周期"的脉络串起来。
盲区清单(速查表)
| # | 盲区 | 关键词 | 出处 |
|---|---|---|---|
| 1 | self._handle_connection 传给 start_server 是回调函数 | client_connected_cb、回调 | S0-Q5 |
| 2 | pass 必须写(try/except 语法要求) | pass、占位 | S0-Q5 |
| 3 | JSON-RPC over TCP 实现进程间通信 | JSON-RPC、socket、IPC | S1-Q28 |
| 4 | EventBus 进程内 vs Kafka/Redis 跨进程 | EventBus、Pub/Sub | S1-Q28 |
| 5 | daemon 里特殊 handler 订阅 EventBus 转 JSON 通过 socket 发 | handler、桥接 | S1-Q28 |
| 6 | 双进程架构 S1→S2 演进 | daemon、client、拆进程 | S1-Q73、S2-Q1 |
| 7 | daemon 常驻 + client 发指令显示 | 常驻服务、瘦客户端 | S1-Q73 |
| 8 | SocketClient 是通用客户端(CLI 和 TUI 都用) | SocketClient、复用 | S2-Q33 |
| 9 | _pending 字典解决异步请求-响应时间差 | Future、配对 | S2-Q5、Q38 |
| 10 | _dispatch 消息分流(JSON-RPC 响应 vs 事件推送) | _dispatch、分流 | S2-Q6、Q26、Q31-32 |
| 11 | send_command + run_event_loop 双协程设计 | 双协程、分流 | S2-Q41、Q47 |
| 12 | 1 TCP + 2 协程 vs 2 TCP + 1 协程方案对比 | 方案对比、复用连接 | S2-Q47 |
| 13 | TUI 纯观察者,只发 event.subscribe 不发 agent.run | 观察者、订阅 | S2-Q65 |
| 14 | daemon 同时只能跑一个 agent run(_current_run_task 限流) | 互斥、限流 | S2-Q22、Q64 |
| 15 | Fan-out 扇出概念 | fan-out、扇出 | S2-Q51 |
| 16 | IPC topic/scope 双重过滤 fnmatch | fnmatch、通配符 | S2-Q50 |
| 17 | EventPushEnvelope kind 字段扩展点 | kind、扩展点 | S2-Q48 |
| 18 | 客户端 exitcode 退出码 + CLI 断开 agent 继续跑 | exitcode、解耦 | S2-Q43-45 |
| 19 | TUI 重连 vs CLI 不重连 | 重连、生命周期 | S2-Q52 |
| 20 | 非阻塞写入 队列 + drain task | 队列、drain、非阻塞 | S3-Q1-2 |
| 21 | SIGKILL 丢数据 vs SIGTERM 不丢 | 信号、优雅关闭 | S3-Q1 |
| 22 | S3 多 Run 并发替代 S2 互斥 | 并发、set 替代 | S3-Q3 |
| 23 | 子进程层级 mini-core/mini-tui vs BashTool 子进程 | 进程层级、subprocess | S3-Q22 |
逐条详解
1. self._handle_connection 传给 start_server 是回调函数
原来怎么理解的:看到 asyncio.start_server(self._handle_connection, ...) 这行,不理解"把一个方法名传进去"是什么意思,以为是要立即执行它,或者它会被循环调用。 正确解释:self._handle_connection 是一个回调函数(callback),对应 asyncio 文档里的 client_connected_cb 参数。start_server 不会立即调用它,而是把它"登记"下来——每当有新的客户端 TCP 连接进来,事件循环就自动调用一次这个回调,把 (reader, writer) 两个参数传给它。程序员不需要写 while True: accept() 循环,框架内部做了 accept → create_task(callback)。 代码示例:
# socket_server.py
async def start(self):
# self._handle_connection 是回调,每来一个连接自动调用一次
self._server = await asyncio.start_server(
self._handle_connection, host, port
)
async def _handle_connection(self, reader, writer):
# 每个客户端连接 = 一个独立协程跑这个函数
while True:
line = await reader.readline()
...
延伸:类比 Java 的 ServerSocket.accept() 需要自己写循环 + 手动 new Thread 处理;asyncio 把这套封装成"传个回调就行"。回调本质是"控制反转"——你把函数交给框架,框架在合适的时机替你调用。
2. pass 必须写(try/except 语法要求)
原来怎么理解的:在 _handle_connection 的 try/except 块里看到 except: pass,疑惑这个 pass 是不是可以省略,或者只是个占位符随手写的。 正确解释:pass 必须写。Python 语法规定 try/except/def/class/if 等"需要有代码体"的结构,体内至少要有一行语句。如果 except 块里什么都不写,解释器会报 IndentationError。pass 是一个空操作语句(no-op),存在的唯一目的就是"占一行",让语法成立。它表示"我故意在这里什么都不做"。 代码示例:
try:
line = await reader.readline()
except Exception:
pass # ← 删掉这行会语法报错;表示"吞掉异常不处理"
延伸:Python 没有 Java 的 {} 块分隔,靠缩进 + 必须有内容。pass 还可用于"先搭骨架"——def foo(): pass 先占位再慢慢实现。另外三个类似语句:...(Ellipsis)、raise、return,都能单独成行满足语法。
3. JSON-RPC over TCP 实现进程间通信
原来怎么理解的:以为进程间通信要么得用 Kafka/Redis 这种中间件,要么得用 Python 的 multiprocessing Pipe,不清楚"裸 TCP + JSON"怎么就成了 IPC。 正确解释:MiniClaude 的 daemon 和 client 是同机两个进程,通过 localhost TCP socket 通信,传输的内容是 JSON-RPC 2.0 格式的文本(一行一个 JSON,以 \n 分隔)。JSON-RPC 是个标准协议:请求是 {jsonrpc, method, params, id},响应是 {jsonrpc, id, result} 或 {jsonrpc, id, error}。不需要中间件,不需要二进制序列化,TCP 只是"管道",JSON-RPC 是"管道里跑的报文格式"。 代码示例:
# 客户端发请求
req = JsonRpcRequest(jsonrpc="2.0", method="agent.run",
params={"goal": "读文件"}, id=req_id)
writer.write(req.model_dump_json().encode() + b"\n")
# 服务端按行读
line = await reader.readline()
msg = json.loads(line) # 解析成 dict
method = msg["method"] # "agent.run"
handler = self._handlers[method] # 字典路由
result = await handler(msg["params"])
延伸:JSON-RPC 类似 Java 里的 HTTP + JSON 调用,但省掉了 HTTP 头开销,适合长连接。本机 TCP 走 loopback,延迟在微秒级,比真网络快得多。一行一 JSON(JSONL)的好处是 readline() 天然分帧,不用自己算长度前缀。
4. EventBus 进程内 vs Kafka/Redis 跨进程
原来怎么理解的:以为 EventBus 就是"小型 Kafka/Redis",应该也能跨进程广播消息。 正确解释:EventBus 是进程内通信,本质就是一个 list[handler] + for + await handler(event)——就是普通函数调用,不跨进程、不跨机器、不走网络。publish(event) 等价于"遍历列表挨个调函数"。Kafka/Redis Pub/Sub 是跨机器跨进程的消息中间件,有 broker、有网络、有持久化。两者是不同层面的东西:EventBus 只管一个进程内部解耦,跨进程由 socket 负责。 代码示例:
class EventBus:
def __init__(self):
self._subs: list[Callable] = []
def subscribe(self, h):
self._subs.append(h) # 就是 list.append
async def publish(self, event):
for h in list(self._subs): # 就是 for 循环调函数
await h(event) # 没有网络、没有队列、没有 broker
延伸:Java 类比——EventBus 像 Guava EventBus(纯进程内),Kafka 像 Kafka(分布式流)。MiniClaude 选 EventBus 是因为 daemon 内部解耦用进程内足够,跨进程才上 socket,避免过度设计。
5. daemon 里特殊 handler 订阅 EventBus 转 JSON 通过 socket 发
原来怎么理解的:不理解"EventBus 是进程内的,那客户端(另一个进程)怎么收到事件"。 正确解释:daemon 里有一个特殊 handler(IpcEventBroadcaster),它做两件事:①订阅 EventBus,跟其他 handler 一样收到进程内所有事件;②收到事件后把它转成 JSON,通过 socket 的 writer.write() 发给订阅了的客户端。它是个"桥接器"——把进程内的函数调用,桥接成跨进程的 TCP 报文。EventBus 自己不知道有 socket,socket 也不知道有 EventBus,桥接器在中间翻译。 代码示例:
# 桥接器:订阅 EventBus → 转 JSON → 发 socket
class IpcEventBroadcaster:
def __init__(self, trace: TraceWriter | None = None):
self._subscriptions: list[_Subscription] = []
self._trace = trace
# 注意:构造器里不订阅 EventBus,订阅在 CoreApp.run() 里做:
# self._broadcaster = IpcEventBroadcaster(trace=self._trace)
# self._bus.subscribe(self._broadcaster.handle) ← 方法名是 handle
async def handle(self, event: BaseModel) -> None:
event_dict = event.model_dump() # event 字段是 dict(已 model_dump)
event_type = event_dict.get("type", "")
run_id = event_dict.get("run_id")
for sub in list(self._subscriptions):
if not self._matches_topic(event_type, sub.topics):
continue
if not self._matches_scope(run_id, sub.scope):
continue
envelope = EventPushEnvelope(event=event_dict) # kind 默认 "event"
sub.writer.write(envelope.model_dump_json().encode() + b"\n")
await sub.writer.drain()
延伸:这是经典的"适配器/桥接"模式——把一个接口(EventBus 的 publish)适配成另一个接口(TCP write)。类比 Java:像把一个 ApplicationEvent 转成 STOMP 帧通过 WebSocket 推给前端。
6. 双进程架构 S1→S2 演进
原来怎么理解的:S1 时以为"AgentRunner 跑在 CLI 进程里"是个临时状态,不清楚 S2 为什么要拆成两个进程。 正确解释:S1:用户敲 mini run,所有代码(AgentRunner + Loop + provider + tool + printer)在同一个 Python 进程内运行,跑完就退出。S2:拆成两个进程——mini-core(daemon)后台常驻跑 Agent 逻辑,mini run 变成客户端只发指令 + 显示结果,两者通过 socket 通信。演进动机:①常驻服务不用每次重启加载;②支持多客户端(CLI + TUI 同时连);③UI 和逻辑解耦。关键:所有 S1 模块接口不动,直接复用,EventBus 保持不变,只加一个"socket 广播 handler"。 代码示例:
# S1: 单进程
mini run --goal "..." → AgentRunner.run() → 跑完退出
# S2: 双进程
mini-core → daemon 常驻,监听 TCP
mini run --goal "..." → client,发 agent.run,订阅事件显示
mini tui → client,纯观察
延伸:类比 Java——S1 像一个 main 方法里跑完所有逻辑;S2 像把后端拆成 Spring Boot 服务 + 前端 HTTP 客户端。这种"先单体后拆分"的演进路径在生产项目里非常常见。
7. daemon 常驻 + client 发指令显示
原来怎么理解的:不清楚"常驻"到底常驻什么,client 发完指令后 agent 在哪跑。 正确解释:daemon 进程一直活着(除非手动 kill),里面跑着事件循环 + EventBus + agent。client(CLI/TUI)是瘦客户端:只负责组装 JSON-RPC 命令通过 TCP 发出去,订阅事件流显示给用户。agent 的实际执行(AgentLoop、调 LLM、调工具)全在 daemon 进程里。client 发完 agent.run 命令后,daemon 用 create_task 把 runner 丢到后台跑,立即返回 run_id,client 拿着 run_id 订阅事件流看进度。 代码示例:
# client 端(瘦)
await client.send_command("agent.run", {"goal": "..."}) # 发指令
await client.run_event_loop() # 收事件显示
# daemon 端(胖,常驻)
async def _handle_agent_run(self, params):
run_task = asyncio.create_task(runner.run(...)) # 丢后台
self._current_run_task = run_task
return {"run_id": run_id} # 立即返回,不等 agent 跑完
延伸:这是"胖服务端 + 瘦客户端"架构。好处是 agent 跑几小时也不受 client 退出影响(见盲区 18)。类比 Java RMI 或 gRPC——客户端只持有 stub,真正干活在服务端。
8. SocketClient 是通用客户端(CLI 和 TUI 都用)
原来怎么理解的:以为 SocketClient 是 TUI 专属的,CLI 另有一套通信代码。 正确解释:SocketClient 是通用客户端类,CLI 和 TUI 都用它。它封装了:发 JSON-RPC 命令(send_command)、收消息分流(_dispatch + run_event_loop)、Future 配对(_pending)。CLI 和 TUI 的区别只在呈现层:CLI 注册的 on_event 回调是 StdoutPrinter(打印到终端),TUI 注册的是 _handle_event(更新 Textual 界面)。通信层完全复用。 代码示例:
# CLI 侧
client.on_event(printer.handle) # StdoutPrinter 打印
# TUI 侧
client.on_event(self._handle_event) # 更新 RichLog widget
# 两者都调同一套:
await client.send_command("event.subscribe", {...})
await client.run_event_loop()
延伸:这是"策略模式"——通信骨架固定,把"收到事件怎么办"作为可替换的回调注入。类比 Java:一个 HttpClient 配不同的 ResponseHandler。
9. _pending 字典解决异步请求-响应时间差
原来怎么理解的:看到 self._pending[req.id] 以为是"等待队列",不清楚它在等什么、怎么被唤醒。 正确解释:_pending 是 dict[str, asyncio.Future],解决的核心问题是**"客户端发命令和收响应是异步的"**。send_command 发请求后不能阻塞整个客户端干等这一条响应(因为同一 TCP 上还有事件推送在流入),所以它:①创建一个 Future 存进 _pending[req.id];②await fut 把自己挂起;③后台的 _dispatch 收到匹配 id 的响应时,fut.set_result(result) 把结果塞进去;④send_command 被唤醒拿到结果。Future 就是"取号条",发完请求拿号等着,叫到你的号(id 匹配)就唤醒。 代码示例:
async def send_command(self, method, params):
req_id = str(uuid.uuid4()) # 用 UUID 作请求 id
req = JsonRpcRequest(id=req_id, method=method, params=params)
fut = asyncio.get_running_loop().create_future()
self._pending[req_id] = fut # 存号
self._writer.write(req.model_dump_json().encode() + b"\n") # 发请求
await self._writer.drain()
return await fut # 挂起等叫号
async def _dispatch(self, line: bytes):
msg = json.loads(line) # _dispatch 只解析分流,不读 socket
if "jsonrpc" in msg:
req_id = msg.get("id")
if req_id and req_id in self._pending:
fut = self._pending.pop(req_id)
if not fut.done():
fut.set_result(msg.get("result") or {}) # 叫号,唤醒 send_command
延伸:Future 是"单写多读"的一次性容器——只能 set_result 一次,所以配 not fut.done() 防重复。类比 Java 的 CompletableFuture:fut.complete(result) 唤醒 fut.get()。
10. _dispatch 消息分流(JSON-RPC 响应 vs 事件推送)
原来怎么理解的:以为客户端收消息就一个通道一个用途,不理解为什么需要"分流"。 正确解释:同一条 TCP 连接上混着两种报文:①JSON-RPC 响应(有 id 字段,对应 send_command 的请求);②事件推送(EventPushEnvelope,有 kind: "event" 字段,是 daemon 主动推的)。run_event_loop 负责持续 readline() 读每一条消息,每读到一行就交给 _dispatch(line) 解析分流。_dispatch 只负责"解析 + 分流"——根据"有没有 jsonrpc 字段"和"有没有 kind 字段"判断走哪条路:有 jsonrpc + id 在 _pending 里 → 唤醒对应的 Future;有 kind=event → 调所有注册的 on_event 回调。 代码示例:
async def run_event_loop(self):
# run_event_loop 负责 readline
while True:
line = await self._reader.readline()
if not line:
break
await self._dispatch(line) # 把行丢给 _dispatch 解析分流
async def _dispatch(self, line: bytes) -> None:
# _dispatch 只负责解析分流,不读 socket
msg = json.loads(line)
# 分支 1:JSON-RPC 响应(配对 send_command)
if "jsonrpc" in msg:
req_id = msg.get("id")
if req_id and req_id in self._pending:
fut = self._pending.pop(req_id)
if not fut.done():
if "error" in msg:
fut.set_exception(IpcError(...))
else:
fut.set_result(msg.get("result") or {})
# 分支 2:事件推送
elif msg.get("kind") == "event":
event_data = msg.get("event", {})
for handler in self._event_handlers: # 调所有注册的回调
await handler(event_data)
延伸:这种"一条通道多种报文"是协议设计的常见取舍——复用连接省资源,但接收方必须能区分报文类型。HTTP/2 的多路复用、WebSocket 的 opcode 都是同样思路。
11. send_command + run_event_loop 双协程设计
原来怎么理解的:不理解一个客户端进程为什么需要两个协程,"两个协程不存在 CPU 竞争"那为什么要分开。 正确解释:两个协程的分工:①send_command(主协程)——发命令、await 等响应;②run_event_loop(create_task 创建的后台 task)——持续 readline() 收消息并分流。需要两个协程的根因不是 CPU 并发,而是消息复用需要分流:同一 TCP 上同时有 RPC 响应和事件推送流入,必须有一个人专门"守着读",另一个人"发命令等结果"。如果只有一个协程,发完命令阻塞等响应期间,事件推送就没人收,会堵在 TCP 缓冲区。 代码示例:
async def _run_async(self):
client = SocketClient(...)
# 后台 task:守着读
asyncio.create_task(client.run_event_loop())
# 主协程:发命令
await client.send_command("event.subscribe", {...})
await client.send_command("agent.run", {"goal": "..."})
延伸:协程之间不是"并行计算"而是"并发 IO"——一个 await 让出时另一个接着跑,单线程内交替。类比 Java:像开一个读线程 + 一个写线程,但 Python 用协程省掉了线程开销。
12. 1 TCP + 2 协程 vs 2 TCP + 1 协程方案对比
原来怎么理解的:觉得"两个协程"挺别扭,疑问"为什么不分开两条 TCP 各管各的"。 正确解释:两种方案都能解决分流问题:
-
当前方案:1 条 TCP + 2 协程。RPC 响应和事件推送共用一条连接,靠
_dispatch在应用层分流。优点:省连接、省端口、客户端只管一个 socket。缺点:应用层要写分流逻辑。 -
备选方案:2 条 TCP + 1 协程。一条专走 RPC(请求-响应),一条专走事件推送。优点:职责清晰不用分流。缺点:多占端口和连接资源、客户端要管两个 socket 的生命周期。
MiniClaude 选前者,因为本机通信连接便宜,但分流逻辑写一次就够,更简单。 代码示例:
# 方案 A(当前):1 TCP + 2 协程
sock = connect()
create_task(read_loop(sock)) # 协程 1 收
await send_command(sock, ...) # 协程 2 发+等
# 方案 B(备选):2 TCP + 1 协程
rpc_sock = connect() # 专走 RPC
evt_sock = connect() # 专走事件
# 一个协程 select 两个 sock
延伸:这是经典的"复用 vs 专用"权衡。HTTP/1.1 用单连接多请求(要分流),HTTP/2 用多路复用(帧分流),gRPC 一个 HTTP/2 连接跑所有 RPC。专用连接(如数据库连接池)则反过来。
13. TUI 纯观察者,只发 event.subscribe 不发 agent.run
原来怎么理解的:以为 TUI 也能像 CLI 一样发 agent.run 启动 agent,或者一个 TUI 对应一个 agent。 正确解释:S2 阶段 TUI 是纯观察者——只发 event.subscribe 订阅事件流,不发 agent.run。agent 由 CLI 启动,TUI 只是"监控面板",订阅的是 daemon 中 EventBus 上的事件。典型用法:开三个终端——一个跑 daemon、一个跑 TUI、一个敲 mini run。TUI 能"看"不能"控制"。
⚠️ 阶段说明:此为 S2 阶段状态。S3+ 之后 TUI 也能通过
session.send_message启动 agent,不再仅是纯观察者。下面的代码示例保留 S2 历史描述。
代码示例:
# CLI: 发 agent.run 启动 agent
await client.send_command("agent.run", {"goal": "..."})
# TUI(S2 阶段): 只订阅,不发 agent.run
await client.send_command("event.subscribe",
{"topics": ["run.*", "tool.*"], "scope": "global"})
# 之后只接收事件,更新界面
延伸:这是"CQRS"思想的一种体现——命令(写)和查询(读)分离。CLI 是命令端,TUI 是查询端。后续阶段 TUI 才可能加发命令的能力。
14. daemon 同时只能跑一个 agent run(_current_run_task 限流)
原来怎么理解的:不清楚多个客户端同时发 agent.run 会怎样,以为能并发跑多个 agent。 正确解释:S2 的 daemon 同时只能跑一个 agent run,靠 _current_run_task 单变量限流。第二个 agent.run 请求来时,检查到 _current_run_task 还没 done,直接抛 RuntimeError("a run is already in progress") 拒绝。所有 TUI 看的都是同一个 agent 的事件流。TUI 不"操作" agent,只能看。
⚠️ 阶段说明:此为 S2 阶段状态。当前源码已演进为 S3 的
_running_runs: set[asyncio.Task[Any]],不再用_current_run_task互斥,改为多 Run 并发 +add_done_callback(discard)自动清理(见盲区 22)。下面的代码示例保留 S2 历史描述。
代码示例:
# S2 互斥限流(仅作历史对照,当前源码已演进为 S3 的 _running_runs 集合)
self._current_run_task: asyncio.Task | None = None
async def _handle_agent_run(self, params):
if (self._current_run_task is not None
and not self._current_run_task.done()):
raise RuntimeError("a run is already in progress") # 拒绝
self._current_run_task = asyncio.create_task(runner.run(...))
return {"run_id": run_id}
延伸:这是"单写者"互斥。S3 改成了多 Run 并发(见盲区 22)。类比 Java 的 synchronized 单例锁,或者"同一时刻只能有一个 job 在跑"的简易调度器。
15. Fan-out 扇出概念
原来怎么理解的:看到笔记里写"fan-out 完成后",以为是某个英文缩写,不懂什么意思。 正确解释:Fan-out 不是缩写,是电子工程术语"扇出"——一个信号源驱动多个负载。在事件系统里:一个事件从 EventBus 出发,扇出到多个订阅者(EventWriter + IpcEventBroadcaster → 多个客户端)。"fan-out 完成"意味着所有订阅者都处理完了这条事件。对应地,多个输入汇聚到一个点叫 fan-in。 代码示例:
# 一个事件 → 扇出到多个订阅者
async def publish(self, event):
for h in list(self._subs): # fan-out:1 个事件 → N 个订阅者
await h(event)
# 订阅者:EventWriter(落盘) + IpcEventBroadcaster(推 client) + ...
延伸:扇出在硬件里指一个门电路输出连多少个后续门,软件里借指"一发多收"。Kafka 的 consumer group、Spring 的 ApplicationEventMulticaster 都是 fan-out 实现。
16. IPC topic/scope 双重过滤 fnmatch
原来怎么理解的:看到 IpcEventBroadcaster 的 topic/scope 过滤,不懂 fnmatch 是什么、为什么要双重过滤。 正确解释:daemon 推事件给客户端时,要按客户端的订阅偏好过滤,避免推一堆客户端不感兴趣的事件。双重过滤:①topic 用 fnmatch 通配符匹配(如订阅 "run.*" 能匹配 "run.started"、"run.finished");②scope 按 run_id 过滤——"global" 全通(推所有事件),"run:<run_id>" 只推该 run 的事件。两层都过才推。 代码示例:
# IpcEventBroadcaster.handle 中
for sub in list(self._subscriptions):
# 第一层:topic 通配符匹配
if not any(fnmatch.fnmatch(event_type, t) for t in sub.topics):
continue
# 第二层:scope 按 run_id 匹配(global 全通,run:<id> 精确匹配)
if not self._matches_scope(run_id, sub.scope):
continue
sub.writer.write(envelope_json + "\n") # 都过了才推
@staticmethod
def _matches_scope(run_id: str | None, scope: str) -> bool:
if scope == "global":
return True
if scope.startswith("run:"):
return run_id == scope[4:] # 去掉 "run:" 前缀后比对 run_id
return False
延伸:fnmatch 是 Python 标准库,语法和 shell glob 一样(* 任意、? 单字符、[seq] 字符集)。类比 Java 的 String.matches() 但更轻。这种"topic 通配 + scope 精确"的双层过滤在消息中间件里很常见(如 MQTT 的 topic 层级 + Redis 的 channel pattern)。
17. EventPushEnvelope kind 字段扩展点
原来怎么理解的:看到 EventPushEnvelope 的 kind 字段,疑惑"除了 event 还能有别的吗"。 正确解释:目前 kind 只有 "event" 一种值。但从设计角度看,kind 是预留的扩展点——将来可能有 kind: "error"(系统级错误推送)、kind: "system"(系统通知)等。_dispatch 根据 kind 分流到不同处理逻辑。现在只用一种不代表字段多余,它是"为未来留口子"。 代码示例:
class EventPushEnvelope(BaseModel):
kind: str = "event" # 目前只有 "event",预留扩展
event: Event | None = None
# _dispatch 分流
if msg.get("kind") == "event":
await self._on_event(msg["event"])
elif msg.get("kind") == "error": # 将来可能
...
延伸:这是"开闭原则"的体现——对扩展开放(加新 kind 不改老逻辑),对修改封闭。类比 HTTP 的状态码类别(1xx/2xx/3xx...),预留了语义空间。
18. 客户端 exitcode 退出码 + CLI 断开 agent 继续跑
原来怎么理解的:看到 exitcode 变量觉得"也没传出去啊",不清楚它干嘛;也不清楚 CLI 进程断了 agent 会不会跟着死。 正确解释:①exitcode 是客户端进程的退出码,通过 asyncio.run(_run_async()) 的返回值传到同步层 cmd_run,再 sys.exit(exitcode)。退出码 1 表示失败(通信失败或 run_event_loop 出错),0 表示正常。区分"哪种失败"靠日志的异常信息,不靠退出码。②CLI 进程断开时,run_event_loop 因通信失败退出,但 agent 在 daemon 进程里继续跑——CLI 只是"看的人走了",不是"把 agent 关了"。TUI 如果还连着,能继续看到事件。 代码示例:
async def _run_async() -> int:
try:
await client.send_command("agent.run", ...)
await client.run_event_loop()
return 0
except IpcError:
return 1 # 通信失败 → 退出码 1
# 同步层
exitcode = asyncio.run(_run_async())
sys.exit(exitcode) # ← 传出去了
# 此时 CLI 死了,但 daemon 里的 agent.run task 还在跑
延伸:这是"控制面与执行面分离"的好处——客户端是控制/观察通道,断了不影响执行面。类比 Kubernetes:kubectl 断了不影响集群里跑着的 Pod。退出码是 Unix 惯例,shell 脚本靠 $? 判断成败。
19. TUI 重连 vs CLI 不重连
原来怎么理解的:以为客户端断线后的行为都一样,不清楚 TUI 和 CLI 在重连上是有意区分的。 正确解释:run_event_loop 结束条件:daemon 断开、daemon 出错、或 agent run 完成后自己 client.close()。CLI 不重连——run_event_loop 退出就退出整个 CLI 进程。TUI 重连——外层 _socket_loop 有 while True,断线后会重新建立 socket 连接,配合 --replay RUN_ID 把断线期间错过的历史事件补回来再看实时流。这是设计上有意为之的区别:CLI 是"一次性任务",TUI 是"持续监控"。 代码示例:
# TUI 的 _socket_loop
async def _socket_loop(self):
while True: # ← TUI 重连的关键
client = SocketClient(self._host, self._port)
try:
await client.connect()
except (ConnectionRefusedError, OSError): # 连不上 daemon
await asyncio.sleep(2) # 等 2 秒再重连
continue
...
await client.run_event_loop()
# 无论正常断开还是异常,外层都会再 sleep(2) 后重试
await asyncio.sleep(2)
# CLI 没有这层 while True,断了就退出
延伸:这是"断线恢复策略"按角色定制。CLI 任务跑完就退,重连没意义;TUI 是长监控,必须自愈。类比 Java:批处理 job 失败就失败,Web 长连接要带指数退避重连。
20. 非阻塞写入 队列 + drain task
原来怎么理解的:看到 TraceWriter 说"非阻塞写入:队列 + drain task",不懂为什么不直接写文件。 正确解释:emit() 不直接写文件,而是把记录扔进 asyncio.Queue 就返回;另一个后台 task(_drain)慢慢从队列取数据写文件。为什么需要:emit() 在 asyncio 事件循环里被调用,如果直接写文件,哪怕只阻塞 5 毫秒,整个事件循环都停了——所有协程一起等。扔队列是微秒级操作,不卡事件循环;写文件的慢活交给后台 drain task。 代码示例:
class TraceWriter:
def __init__(self):
self._queue: asyncio.Queue = asyncio.Queue()
self._drain_task = asyncio.create_task(self._drain())
def emit(self, record): # 同步方法,非 async
self._queue.put_nowait(record) # 扔进去立刻返回
async def _drain(self): # 后台慢慢写
while True:
record = await self._queue.get()
await self._write_file(record)
self._queue.task_done()
async def stop(self):
await self._queue.join() # 等队列清空
self._drain_task.cancel()
延伸:Java 类比 Log4j 的 AsyncAppender——生产者不阻塞,消费者异步写。这是"生产者-消费者"模式解耦"产生速度"和"消费速度"。
21. SIGKILL 丢数据 vs SIGTERM 不丢
原来怎么理解的:不清楚为什么进程被杀时数据会不会丢,SIGKILL 和 SIGTERM 有什么区别。 正确解释:SIGKILL(kill -9)是操作系统强制立即终止进程,进程没机会执行任何清理代码——TraceWriter 队列里还没写出的记录会丢失。SIGTERM(kill,默认信号)是"礼貌地请进程退出",进程能收到信号并执行关闭流程:TraceWriter 的 stop() 会 await queue.join() 等队列清空再退出,不丢数据。所以优雅关闭必须用 SIGTERM,绝不能用 SIGKILL 除非进程卡死。 代码示例:
# 优雅关闭流程(响应 SIGTERM)
async def shutdown(self):
await self._trace.stop() # ← 内部 await queue.join() 等队列清空
# SIGKILL 不会走到这里,队列里没写出的记录直接丢
# MiniClaude 用 loop.add_signal_handler() 把 SIGTERM 接进事件循环
loop.add_signal_handler(signal.SIGTERM,
lambda: asyncio.create_task(shutdown()))
延伸:SIGKILL 不可被捕获、不可被忽略,是"最后手段"。类比 Java 的 Thread.interrupt() vs kill -9——前者礼貌通知能响应,后者直接枪毙。生产环境 graceful shutdown 都是先 SIGTERM 等几秒再 SIGKILL 兜底。
22. S3 多 Run 并发替代 S2 互斥
原来怎么理解的:记得 S2 说"同时只能跑一个 agent run",不清楚 S3 改成什么了。 正确解释:S3 把 S2 的"单 Run 互斥"改成了"多 Run 并发"。对比:S2 用单变量 _current_run_task,第二个请求直接 raise RuntimeError 拒绝。S3 用 set[asyncio.Task[Any]] 集合 _running_runs,每个 agent.run 都 create_task 加入集合,add_done_callback(discard) 完成后自动从集合移除,不再拒绝。为什么 S3 要支持并发:S3 引入了 Task 系统,用户可能同时创建多个任务,每个任务触发独立的 Agent Run。注意:并发的是 Run(Agent 执行实例),不是 LLM 调用——每个 Run 有自己的 AgentLoop、ExecutionContext、TaskManager,互不干扰。 代码示例:
from typing import Any
# S2:互斥(仅作历史对照,当前源码已演进为 S3)
self._current_run_task: asyncio.Task | None = None
if self._current_run_task is not None and not self._current_run_task.done():
raise RuntimeError("a run is already in progress") # 拒绝
# S3:并发
self._running_runs: set[asyncio.Task[Any]] = set()
# runner 在 SessionManager 内部创建,handler 只调 send_message
run_task = asyncio.create_task(
self._sessions.send_message(session.id, cmd.goal, run_id=run_id)
)
self._running_runs.add(run_task)
run_task.add_done_callback(self._running_runs.discard) # 完成自动清理
return AgentRunResult(run_id=run_id) # 立即返回,不拒绝
延伸:set + discard 是协程安全的标准清理套路。discard 比 remove 安全(元素不存在不报错)。关闭时批量 cancel() 集合里所有 task。类比 Java:ExecutorService + Future 集合,shutdown 时 shutdownNow() 批量中断。
23. 子进程层级 mini-core/mini-tui vs BashTool 子进程
原来怎么理解的:以为 BashTool 起的子进程和 mini-core、mini-tui 是同一层级的进程。 正确解释:它们是不同层级的进程。mini-core(daemon)和 mini-tui 是平级的独立进程,由用户在终端分别启动,通过 TCP 通信,互不管理。而 BashTool 创建的子进程是 daemon 的子进程——daemon 进程内通过 asyncio.create_subprocess_shell("ls -la") 启动,PID 不同、生命周期短暂(跑完就退出)、继承 daemon 的权限和环境。两者关系完全不同:前者是"兄弟",后者是"父子"。 代码示例:
用户启动终端
├── mini-tui 进程(独立进程,平级)──┐
└── mini-core 进程(独立进程,平级)─┤ TCP 通信
└── asyncio.create_subprocess_shell("ls -la")
→ 临时子进程(daemon 的子进程,跑完退出)
| <br /> | mini-core / mini-tui | BashTool 的子进程 |
|---|---|---|
| 生命周期 | 长期运行 | 临时创建,跑完退出 |
| 关系 | 独立进程,TCP 通信 | daemon 的子进程,继承权限 |
| 管理 | 互不管理 | 父进程可以 kill、等待 |
延伸:进程有父子关系(父死子变孤儿),线程全部平级没有"父线程",协程/Task 逻辑上有创建者但机制上无层级(创建者不持有引用就是"孤儿 Task")。BashTool 子进程超时会被 daemon kill,而 daemon 不能 kill mini-tui(只能等它自己退)。
复习自检
-
start_server的第一个参数是回调还是立即执行?谁在调用它? -
pass能不能删掉?为什么? -
EventBus 跨进程吗?跨进程通信靠谁?
-
daemon 里的"桥接 handler"做哪两件事?
-
S1 到 S2 进程数从几个变几个?agent 在哪个进程跑?
-
CLI 和 TUI 共用哪个客户端类?区别在哪一层?
-
_pending存的是什么?send_command怎么被唤醒的? -
_dispatch为什么需要分流?分哪两种报文? -
为什么需要两个协程?根因是 CPU 竞争吗?
-
1 TCP+2 协程 vs 2 TCP+1 协程,各自优缺点?
-
TUI 能发
agent.run吗?它发什么命令? -
S2 daemon 同时能跑几个 agent run?靠什么限流?
-
fan-out 是缩写吗?什么意思?
-
IPC 过滤有哪两层?topic 用什么匹配?
-
EventPushEnvelope.kind目前几个值?为什么留这个字段? -
CLI 断了 agent 会死吗?为什么?
-
TUI 断线重连,CLI 为什么不重连?
-
TraceWriter 为什么不直接写文件?
-
SIGKILL 和 SIGTERM 谁会丢数据?为什么?
-
S3 怎么把 S2 的互斥改成并发?用什么数据结构?
-
BashTool 子进程和 mini-core 是同一层级吗?
易错点总结
-
回调不是立即调用:
start_server(self._handle_connection)是"登记",框架在连接到来时才调。 -
EventBus 不跨进程:它是进程内的函数调用,跨进程靠 socket + JSON-RPC,靠桥接 handler 翻译。
-
双协程根因是分流不是并发:1 条 TCP 混两种报文必须分流,不是为了 CPU 并行。
-
_pending配对的是请求-响应,不是事件推送:事件推送走on_event回调,不走 Future。 -
CLI 断了 agent 不死:agent 在 daemon 里跑,client 只是观察者,别把"客户端退出"当成"任务取消"。
-
SIGTERM 能优雅关闭,SIGKILL 不能:生产环境永远先 SIGTERM。
-
BashTool 子进程 ≠ mini-core/mini-tui:前者是 daemon 的子进程,后两者是平级独立进程。
-
S2 互斥、S3 并发:别把两个阶段的并发模型搞混。
-
SocketClient 是通用的:CLI 和 TUI 都用它,区别只在
on_event回调,别以为 TUI 有专属客户端。 -
kind字段不是只有 event 才合理:预留扩展点是开闭原则的体现,不是冗余。