QZ Site
HomeBlogProjectsAbout

© 2026 QZ Site. All rights reserved.

豫ICP备2026034998号

← Back to blog

05-IPC 与网络通信盲区梳理

MiniClaude·August 1, 2026·81 min read

IPC 与网络通信盲区梳理

概述

本篇梳理在学习 MiniClaude 项目过程中,围绕进程间通信(IPC)、TCP 网络通信、双进程架构、消息分流暴露的知识盲区。每条盲区包含:①原来的困惑/错误理解 ②正确解释 ③代码示例 ④延伸知识。

学习者在 Java 阶段习惯了"线程 + RPC 框架"的并发模型,进入 Python asyncio + 裸 TCP + JSON-RPC 的世界后,对"回调函数到底谁在调""为什么一条 TCP 上要跑两个协程""Future 怎么把请求和响应配对"等机制层问题反复卡壳。本篇把这些盲区按"基础回调 → 跨进程架构 → 消息分流 → 并发与生命周期"的脉络串起来。

盲区清单(速查表)

#盲区关键词出处
1self._handle_connection 传给 start_server 是回调函数client_connected_cb、回调S0-Q5
2pass 必须写(try/except 语法要求)pass、占位S0-Q5
3JSON-RPC over TCP 实现进程间通信JSON-RPC、socket、IPCS1-Q28
4EventBus 进程内 vs Kafka/Redis 跨进程EventBus、Pub/SubS1-Q28
5daemon 里特殊 handler 订阅 EventBus 转 JSON 通过 socket 发handler、桥接S1-Q28
6双进程架构 S1→S2 演进daemon、client、拆进程S1-Q73、S2-Q1
7daemon 常驻 + client 发指令显示常驻服务、瘦客户端S1-Q73
8SocketClient 是通用客户端(CLI 和 TUI 都用)SocketClient、复用S2-Q33
9_pending 字典解决异步请求-响应时间差Future、配对S2-Q5、Q38
10_dispatch 消息分流(JSON-RPC 响应 vs 事件推送)_dispatch、分流S2-Q6、Q26、Q31-32
11send_command + run_event_loop 双协程设计双协程、分流S2-Q41、Q47
121 TCP + 2 协程 vs 2 TCP + 1 协程方案对比方案对比、复用连接S2-Q47
13TUI 纯观察者,只发 event.subscribe 不发 agent.run观察者、订阅S2-Q65
14daemon 同时只能跑一个 agent run(_current_run_task 限流)互斥、限流S2-Q22、Q64
15Fan-out 扇出概念fan-out、扇出S2-Q51
16IPC topic/scope 双重过滤 fnmatchfnmatch、通配符S2-Q50
17EventPushEnvelope kind 字段扩展点kind、扩展点S2-Q48
18客户端 exitcode 退出码 + CLI 断开 agent 继续跑exitcode、解耦S2-Q43-45
19TUI 重连 vs CLI 不重连重连、生命周期S2-Q52
20非阻塞写入 队列 + drain task队列、drain、非阻塞S3-Q1-2
21SIGKILL 丢数据 vs SIGTERM 不丢信号、优雅关闭S3-Q1
22S3 多 Run 并发替代 S2 互斥并发、set 替代S3-Q3
23子进程层级 mini-core/mini-tui vs BashTool 子进程进程层级、subprocessS3-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-tuiBashTool 的子进程
生命周期长期运行临时创建,跑完退出
关系独立进程,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 才合理:预留扩展点是开闭原则的体现,不是冗余。

Contents

  • 概述
  • 盲区清单(速查表)
  • 逐条详解
  • 1. self._handle_connection 传给 start_server 是回调函数
  • 2. pass 必须写(try/except 语法要求)
  • 3. JSON-RPC over TCP 实现进程间通信
  • 4. EventBus 进程内 vs Kafka/Redis 跨进程
  • 5. daemon 里特殊 handler 订阅 EventBus 转 JSON 通过 socket 发
  • 6. 双进程架构 S1→S2 演进
  • 7. daemon 常驻 + client 发指令显示
  • 8. SocketClient 是通用客户端(CLI 和 TUI 都用)
  • 9. _pending 字典解决异步请求-响应时间差
  • 10. _dispatch 消息分流(JSON-RPC 响应 vs 事件推送)
  • 11. send_command + run_event_loop 双协程设计
  • 12. 1 TCP + 2 协程 vs 2 TCP + 1 协程方案对比
  • 13. TUI 纯观察者,只发 event.subscribe 不发 agent.run
  • 14. daemon 同时只能跑一个 agent run(_current_run_task 限流)
  • 15. Fan-out 扇出概念
  • 16. IPC topic/scope 双重过滤 fnmatch
  • 17. EventPushEnvelope kind 字段扩展点
  • 18. 客户端 exitcode 退出码 + CLI 断开 agent 继续跑
  • 19. TUI 重连 vs CLI 不重连
  • 20. 非阻塞写入 队列 + drain task
  • 21. SIGKILL 丢数据 vs SIGTERM 不丢
  • 22. S3 多 Run 并发替代 S2 互斥
  • 23. 子进程层级 mini-core/mini-tui vs BashTool 子进程
  • 复习自检
  • 易错点总结