pi 桥改为工作进程池 + homeagent 落盘投递账本

## pi 桥:模型工作下到子进程(并发模型重构)

主进程原来自己跑模型,而 pi 的会话装载是同步的:`SessionManager.open()` 走
`openSync` + `readSync` 循环把整个 `.jsonl` 读进内存并逐行 JSON.parse。实测本机
最大那条会话 23MB,`open` 一次**阻塞事件循环 118ms**;模型跑起来后 SDK 内部还有
更多同步工作。SSE 读循环在那期间完全停住 → 后续邮件卡在 TCP 缓冲区 → 久到
Gateway 认为连接死了 → 重连 → 重放。

上一轮我在几个调用点前加 `setImmediate` 是无效的仪式(让出一次之后同步工作照样
占满线程),已回退。这一轮把模型工作整体搬进子进程:实测同样的活在 fork 出的
子进程里跑,主进程事件循环阻塞 **0ms**。

- 新增 `src/worker.mjs`:一封邮件一个进程,跑完就退。权限询问期间的挂起只影响
  那一个 worker(原来 `await new Promise(...)` 等人决策,整座桥不再收信)。
- 新增 `src/pool.mjs`:**不同会话并发**(上限 3,每个 worker 约 140MB RSS)、
  **同一会话严格串行**(pi 假定「一文件一持有者」,两个进程同时装载同一条会话
  文件会让各自的内存索引看不见对方追加的行 → 会话树分叉)、满载排队不丢邮件、
  硬超时 SIGKILL 回收卡死进程。
- `src/index.mjs` 只剩 I/O 与调度:SSE、心跳、去重、分派。
- 选进程而不是 `worker_threads`:模型会跑 bash/write/edit,一次 OOM 不该带走
  整座桥。两者实测都能建起 AgentSession,但线程与主线程共享堆和生命周期。
- 「接管会话短暂持有」那套机制(adopted / adoptTimers / releaseAdopted + 兜底
  计时器)整个删掉 —— worker 退出**就是**释放,且普通会话与接管会话一视同仁。
- 轮次超时 60s → 10 分钟:60s 那个数字是「主进程要腾出手收下一封」的产物,
  worker 没有这个理由,等真结论更准(带工具调用的一轮跑几分钟很正常)。
- IPC 只传路径与标量(sessionFile / cwd / grants / 命名指纹)—— AgentSession
  跨不了进程边界,worker 每次从 sessionFile 重新装载。

`test/pool.test.mjs` +19 例,真 fork 子进程、用桩 worker(不装 SDK)跑毫秒级:
并发上限、同会话串行、不同会话真并发(判据是两个进程的心跳交错,不是 running
map 里有两个条目)、sessionFile/grants/命名指纹跨 worker 传递、config() 每次重取、
硬超时回收、权限决策路由、决策原文透传、worker 退出后清路由、mailDrivenIDs、
kind 透传、stop 先发 shutdown 再杀。跑过三组负向对照确认用例真能抓回归:
拆掉串行守卫 / 不传 sessionFile+grants / 硬超时不杀,对应用例分别失败。

## homeagent:投递去重必须落盘

用户报的重复投递不是上一轮那个 bug。两段提示词的措辞差异指出了来源:
SSE 那段写「你把本轮工作做完」,补投那段写「你把结论说出来就行」。

`deliveredMails` 是进程内的 map,而 homeagent 的插件跑在**子进程**里:

  1. 18:59:38 邮件落库,旧插件进程的 SSE 收到,注入第一次
  2. 同一秒 homed 被重启,那一轮被掐断(`context canceled`)
  3. 18:59:45 新进程起来,`deliveredMails` 是空的
  4. 心跳报 `pending_mails: 1`(第一轮没跑完 → read_inbox 没执行 → 仍未读)
     → catchUp 注入第二次

**不能只记「投过没有」**:那会把「重复」换成「丢件」—— 第 2 步里发件人没收到
回信,而记录说「已投过」→ 永远跳过。丢件比重复严重,重复至少人能看出来。

新增 `ledger.go`:JSONL 账本记两个状态。`completed` 才跳过;`delivered` 但未
`completed` 的仍然重投,但提示词前面插一段说明「上一轮被中断,别把同一件事做
两次」。落在 SDK 的 `Settings().DataDir()`;拿不到时退回 key 文件目录;目录不可
写时退化为纯内存(不比修复前差,也不该让插件起不来)。

- 判定与记录在同一把锁里:SSE 与 catchUp 两个 goroutine 的竞态
- 每行写完 fsync:这个文件的全部意义就是「进程死了之后还算数」
- 坏行跳过而不是报错退出(崩溃时最后一行可能写残)→ 那封退化为重投,安全
- 14 天保留期;过期过半时「临时文件 + rename」压实
- `shortID()` 替代 `id[:8]`:日志不该有能力 panic 掉投递协程

`ledger_test.go` +14 例,含两组负向对照(只记「投过」→ 丢件用例失败;不读账本
→ 跨进程用例失败)。

## 契约文档

`B-7.7`(MUST):子进程形式的插件去重必须落盘且区分「投过」与「跑完」,含事故
时序、两状态表、何时标 completed。已知取舍那节标注投递账本是唯一必须落盘的状态。
验收清单加「模型跑到一半重启宿主」一项。

## 生产验证

- pi 三封 → 三条会话:三个 worker PID 并存,回信「收到 1/2/3」各落自己线索
- pi 同一会话两封:严格串行(收到A 19:26:39 → 收到B 19:26:48,全程单 worker)
- pi 主进程事件循环阻塞 1ms(旧版单进程 open 23MB 一次就 118ms)
- homeagent 正常一封:账本 `c:false` → `c:true`,一封回信
- homeagent 处理中重启:日志「上一轮被中断,带说明重投」,**只有一封 Re:**
- homeagent 再次重启:账本 2 条 completed,不再投递,会话邮件数不变
- gateway 7 包 / web 176+26 / pi 269 / dsh 219 / opencode 201 / homeagent 14
This commit is contained in:
2026-09-04 20:33:44 +08:00
parent c297468819
commit 8e501f041e
8 changed files with 2117 additions and 704 deletions

View File

@ -0,0 +1,277 @@
/**
* 投递工作进程池 —— 主进程侧的调度逻辑。
*
* # 它解决的问题
*
* 桥的主进程唯一的实时职责是读 SSE。模型工作放在主进程里跑会占满事件循环
* pi 的会话装载是同步的23MB 的会话文件 `SessionManager.open` 一次阻塞
* 118ms实测模型跑起来之后 SDK 内部还有大量同步工作SSE 读循环停住,
* 后续邮件卡在 TCP 缓冲区,久到 Gateway 认为连接死了 → 重连 → 重放。
*
* 所以:**收到事件就派给一个子进程,主进程立刻回去读 SSE。**
*
* # 并发与串行的边界
*
* - **不同邮件会话并发**,上限 `maxWorkers`(默认 3。上限的理由是内存
* (每个 worker 约 140MB RSS实测和对上游 provider 的并发请求数。
* - **同一邮件会话串行**。这是正确性要求不是限流pi 没有任何锁机制,
* 它假定「一个文件一个持有者」。两个 worker 同时装载同一条会话文件,各自的
* 内存索引都看不见对方追加的行,算出的 parentId 指向对方不知道的 entry
* → 会话树分叉。串行还顺带保证了同一条线索里两封邮件的先后顺序。
*
* # 排队而不是拒绝
*
* 满载时邮件进 `queue`,有 worker 空出来就派。丢掉邮件是不可接受的:
* 发件人只会看到信发出去后再无音讯。队列无上限 —— 有上限就得决定丢哪封,
* 而任何丢弃策略都比「慢一点」糟。
*
* # 主进程持有什么
*
* 只有**路径与标量**sessionFile / cwd / piSessionId / 「一直同意」表 /
* 命名同步指纹。AgentSession 对象跨不了进程边界worker 每次从 sessionFile
* 重新装载 —— 拿到的是包含 TUI 期间写入的全部历史(这也让「短暂持有」
* 从一套需要计时器兜底的机制退化成「worker 退出就是释放」)。
*
* # IPC 协议
*
* 主进程 → worker
* `{type:'job', kind, data, session:{sessionFile,cwd}, grants, lastSyncedName, config}`
* `{type:'permission_decision', relayKey, decision}`
* `{type:'shutdown'}`
* worker → 主进程:
* `{type:'ready'}` 进程起来了,可以派活
* `{type:'log', line}` 日志(主进程加 pid 前缀)
* `{type:'session_opened', piSessionId, sessionFile, cwd, reused}`
* `{type:'permission_pending', relayKey}` 主进程记下路由表
* `{type:'permission_grant', toolName}` 「一直同意」要跨 worker 活下来
* `{type:'name_synced', signature}` 命名指纹,防下一个 worker 重复 sync
* `{type:'reconfigure', url, agentKey}` connect_to_server 换了坐标
* `{type:'done', ok, error}` 这封处理完了
*/
import { fork } from 'node:child_process';
import { fileURLToPath } from 'node:url';
const WORKER_PATH = fileURLToPath(new URL('./worker.mjs', import.meta.url));
/**
* @param {object} deps
* @param {(...a: any[]) => void} deps.log
* @param {() => object} deps.config 每次派活时取一次allowedModels 会随心跳变)
* @param {(url: string, key: string) => void} deps.onReconfigure
* @param {number} [deps.maxWorkers]
* @param {number} [deps.workerMaxMs] worker 硬超时:卡死的进程必须能被回收
* @param {string} [deps.workerPath] 只为测试存在:换成不装 pi SDK 的桩 worker
* 让调度不变量(并发上限、同会话串行、硬超时)能在毫秒级验证。
*/
export function createWorkerPool({
log, config, onReconfigure,
maxWorkers = 3, workerMaxMs = 600_000, workerPath = WORKER_PATH,
}) {
/** 正在跑的 workermailSessionKey -> {child, mailID, startedAt, timer} */
const running = new Map();
/** 等着派的活,先进先出。 */
const queue = [];
/** relay_key -> mailSessionKey把决策路由回发起询问的那个 worker。 */
const permissionRoutes = new Map();
/**
* 跨 worker 存活的会话状态mailSessionKey -> {sessionFile, cwd, piSessionId,
* grants:Set, lastSyncedName}。
*
* 这是 worker 一封一进程之后仍需在主进程留存的全部东西 —— 下一封邮件靠
* sessionFile 接着谈,靠 grants 不重复问已经「一直同意」过的工具。
*/
const sessionState = new Map();
/**
* 被模型降级换掉的旧 pi 会话 id。
*
* 仍要计入 mail_driven它们已经参与过邮件往来而磁盘上的会话文件
* 不会因为换模型而消失 —— 心跳快照仍会上报它们。
*/
const retired = new Set();
let stopped = false;
/**
* 邮件会话 id 作为串行化的键。
*
* 没有 session_id 的事件(理论上不该有)退回 mail_id那样每封各占一个
* worker不会串行 —— 但它们本来也不属于同一条会话。
*/
const keyOf = (data) => data?.session_id || `mail:${data?.mail_id || Math.random()}`;
function submit(kind, data) {
if (stopped) return;
queue.push({ kind, data, key: keyOf(data) });
pump();
}
function pump() {
if (stopped) return;
for (let i = 0; i < queue.length; i++) {
const job = queue[i];
// 同一会话已有 worker 在跑 → 跳过它,看后面有没有别的会话可以先跑。
// 不能 break那会让一条慢会话把所有别的会话都堵住正是要修的病
if (running.has(job.key)) continue;
if (running.size >= maxWorkers) return;
queue.splice(i, 1);
i--;
spawn(job);
}
}
function spawn(job) {
const state = sessionState.get(job.key) || { grants: new Set(), lastSyncedName: '' };
const child = fork(workerPath, [], {
// stdio 继承worker 里 pi SDK 自己打的东西直接进 journalctl。
// 'ipc' 必须显式列出,否则 process.send 不存在。
stdio: ['ignore', 'inherit', 'inherit', 'ipc'],
});
// 硬超时worker 卡死(模型不返回、权限等不到决策而主进程也没收到事件)
// 时必须能回收,否则那条会话的后续邮件永远排队。
const timer = setTimeout(() => {
log(`worker ${child.pid} 处理 ${job.data?.mail_id} 超过 ${workerMaxMs / 1000}s强杀`);
try { child.kill('SIGKILL'); } catch { /* 已经死了 */ }
}, workerMaxMs);
if (typeof timer.unref === 'function') timer.unref();
const entry = { child, mailID: job.data?.mail_id || '', key: job.key, startedAt: Date.now(), timer };
running.set(job.key, entry);
child.on('message', (msg) => onWorkerMessage(entry, msg));
child.on('exit', (code, signal) => {
clearTimeout(timer);
running.delete(job.key);
for (const [rk, k] of permissionRoutes) if (k === job.key) permissionRoutes.delete(rk);
if (code !== 0) {
log(`worker ${child.pid}mail ${entry.mailID})异常退出 code=${code} signal=${signal || '-'}`);
}
pump();
});
child.on('error', (e) => log(`worker ${child.pid} 出错: ${e?.message || e}`));
// 等 worker 说 ready 再派活fork 返回时子进程的 import 还没跑完,
// 此时 send 的消息会排在 IPC 队列里(能收到,但 ready 让顺序确定)。
child.once('message', function first(msg) {
if (msg?.type !== 'ready') return;
child.send({
type: 'job',
kind: job.kind,
data: job.data,
session: {
sessionFile: state.sessionFile || '',
cwd: state.cwd || '',
},
grants: [...state.grants],
lastSyncedName: state.lastSyncedName || '',
config: config(),
});
});
}
function onWorkerMessage(entry, msg) {
const state = sessionState.get(entry.key) || { grants: new Set(), lastSyncedName: '' };
switch (msg?.type) {
case 'log':
log(`[w${entry.child.pid}] ${msg.line}`);
return;
case 'session_opened':
// 一条会话可能先后用过多个 pi 会话 id模型降级会换会话
// 旧 id 仍计入 mail_driven理由见 retired 的注释。
if (state.piSessionId && state.piSessionId !== msg.piSessionId) {
retired.add(state.piSessionId);
}
state.piSessionId = msg.piSessionId;
state.sessionFile = msg.sessionFile;
state.cwd = msg.cwd;
sessionState.set(entry.key, state);
return;
case 'permission_pending':
permissionRoutes.set(msg.relayKey, entry.key);
return;
case 'permission_grant':
// 「一直同意」必须跨 worker 活着worker 一封一进程,不存的话下一封
// 邮件又问一遍,那个选项就是在骗人。
state.grants.add(msg.toolName);
sessionState.set(entry.key, state);
return;
case 'name_synced':
state.lastSyncedName = msg.signature;
sessionState.set(entry.key, state);
return;
case 'reconfigure':
onReconfigure?.(msg.url, msg.agentKey);
return;
case 'done':
if (!msg.ok) log(`投递 ${entry.mailID} 失败: ${msg.error}`);
return;
default:
return;
}
}
/**
* 把权限决策路由到发起询问的那个 worker。
*
* @returns {boolean} 有没有找到对应的 worker。找不到说明那个 worker 已经退了
* (桥重启、硬超时被杀、或者处理已经结束)—— 调用方据此走 B-4.2 的
* 降级路径(把决策当一封通知投进原会话)。
*/
function routePermission(relayKey, decision) {
const key = permissionRoutes.get(relayKey);
if (!key) return false;
const entry = running.get(key);
if (!entry) {
permissionRoutes.delete(relayKey);
return false;
}
permissionRoutes.delete(relayKey);
entry.child.send({ type: 'permission_decision', relayKey, decision });
return true;
}
/** 这条邮件会话有 worker 在跑吗B-4.2 判断降级路径用)。 */
const hasSession = (mailSessionID) => sessionState.has(mailSessionID);
/**
* 邮件驱动过的 pi 会话 id喂给心跳快照的 `mail_driven` 标记。
*
* 不随 worker 退出而清worker 退了不代表那条会话不再参与邮件往来 ——
* 下一封邮件还会接着谈,而人在补全里需要看到它带着这个标记。
* 重启丢是已知取舍(契约第六节)。
*/
const mailDrivenIDs = () => {
const out = new Set(retired);
for (const st of sessionState.values()) {
if (st.piSessionId) out.add(st.piSessionId);
}
return out;
};
function stop() {
stopped = true;
queue.length = 0;
for (const { child, timer } of running.values()) {
clearTimeout(timer);
// 先 shutdown 让 worker 把未决权限 fail closedB-9.2),再给它一点
// 时间自己退。不直接 SIGKILL那样 pi 侧的 await 不会返回,而 worker
// 里可能正握着会话文件。
try { child.send({ type: 'shutdown' }); } catch { /* 通道已断 */ }
setTimeout(() => { try { child.kill('SIGKILL'); } catch { /* 已经死了 */ } }, 2000).unref?.();
}
}
/** 观测用:现在跑着几个、排了几个。 */
const stats = () => ({
running: running.size,
queued: queue.length,
sessions: sessionState.size,
workers: [...running.values()].map((e) => ({
pid: e.child.pid, mailID: e.mailID, ageMs: Date.now() - e.startedAt,
})),
});
return { submit, routePermission, hasSession, mailDrivenIDs, stop, stats };
}