/** * 投递工作进程池 —— 主进程侧的调度逻辑。 * * # 它解决的问题 * * 桥的主进程唯一的实时职责是读 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}` 这封处理完了 * * # 内存边界 * * `sessionState` 与 `retired` 是**跨 worker 长期存活**的两张表,键来自邮件会话流 * —— 会话数随时间单调增长。两条出口:`forget()`(会话归档,确定性)与 * `BoundedMap`/`BoundedSet` 的上限淘汰(兜底)。缺了它们这里就是常驻进程里 * 一处只增不减的结构。 */ import { fork } from "node:child_process"; import { fileURLToPath } from "node:url"; import { BoundedMap, BoundedSet, MAX_TRACKED_SESSIONS, } from "../lib/bounded.js"; 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 {number} [deps.maxAttempts] 同一封邮件的最大尝试次数(含首次)。 * worker 未回报 `done` 就退出(崩溃、SIGKILL、OOM)时按 1s/2s/… 有界重投; * 超过上限就放弃并留日志 —— 无界重投会把一封必定失败的邮件变成永久活锁。 * @param {string} [deps.workerPath] 只为测试存在:换成不装 pi SDK 的桩 worker, * 让调度不变量(并发上限、同会话串行、硬超时)能在毫秒级验证。 */ export function createWorkerPool({ log, config, onReconfigure, maxWorkers = 3, workerMaxMs = 600_000, maxAttempts = 3, workerPath = WORKER_PATH, }) { /** 正在跑的 worker:mailSessionKey -> {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 BoundedMap(MAX_TRACKED_SESSIONS); /** * 被模型降级换掉的旧 pi 会话 id。 * * 仍要计入 mail_driven:它们已经参与过邮件往来,而磁盘上的会话文件 * 不会因为换模型而消失 —— 心跳快照仍会上报它们。 */ const retired = new BoundedSet(MAX_TRACKED_SESSIONS); 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, attempt = 1) { if (stopped) return; queue.push({ kind, data, key: keyOf(data), attempt }); 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, settled: false, }; 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); // 没收到 `done` 就退出 = 这封邮件**从未处理完**。 // // 这是生产上真实存在的静默丢信路径:worker 被 SIGKILL(硬超时)、 // OOM、或自己崩溃时,`done` 永远不会到达,主进程只看到 exit code。 // 原来这里只记一行日志就 pump() —— 发件人看到信发出去了, // 而那条会话再也不会有人回。 // // 重投而不是直接由主进程回信:worker 崩溃可能是内存/上游瞬时故障, // 重启一个进程真能跑通。有界(maxAttempts)是因为「必定失败」的邮件 // 无界重投会变成永久活锁,而日志里只有一行看不出是同一封在原地打转。 if (!entry.settled && !stopped) { const attempt = job.attempt || 1; if (attempt < maxAttempts) { const delay = attempt * 1000; log( `worker ${child.pid}(mail ${entry.mailID})未回报 done 就退出` + `(code=${code} signal=${signal || "-"}),${delay / 1000}s 后` + `第 ${attempt + 1}/${maxAttempts} 次重投`, ); const retry = setTimeout(() => { if (stopped) return; queue.push({ ...job, attempt: attempt + 1 }); pump(); }, delay); if (typeof retry.unref === "function") retry.unref(); // 退避期间不 pump:否则同一会话会被立刻重投,退避形同虚设 return; } log( `worker ${child.pid}(mail ${entry.mailID})重投 ${maxAttempts} 次仍未完成,放弃` + `(code=${code} signal=${signal || "-"})`, ); } else 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": // 标记「这封真的处理完了」:exit 处理器据此区分「正常收尾」 // 与「未回报就崩溃」(后者要重投)。 entry.settled = true; 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); /** * 忘掉一条已归档会话的全部状态。 * * 归档是个**确定性的终点**:归档后那条会话不可寻址(别名 404),也不会再有 * 新邮件投进来。把它的 sessionState 留着只是占内存,而上限淘汰是「猜」—— * 能确切知道该删的时候就不该依赖猜。 * * 正在跑的 worker **不杀**:归档不是中止指令,模型可能正在写文件;它自己跑完 * 就退,只是那一轮的回信会因为会话已归档而被服务端拦下。 * * @param {string} mailSessionID * @returns {boolean} 是否真的删掉了东西 */ function forget(mailSessionID) { if (!mailSessionID) return false; // peek 而不是 get:这是清理路径,不该把即将删掉的条目刷成「最近活跃」。 const state = sessionState.peek(mailSessionID); // 已归档会话的 pi 会话 id 也不必再报 mail_driven:那个标记的用途是让人在 // 补全里看到「这条在跑邮件」,而已归档的会话不在补全候选里。 if (state?.piSessionId) retired.delete(state.piSessionId); return sessionState.delete(mailSessionID); } /** * 邮件驱动过的 pi 会话 id,喂给心跳快照的 `mail_driven` 标记。 * * 不随 worker 退出而清:worker 退了不代表那条会话不再参与邮件往来 —— * 下一封邮件还会接着谈,而人在补全里需要看到它带着这个标记。 * 重启丢是已知取舍(契约第六节);确定性的清理时机是归档(见 forget)。 * * 返回普通 Set 而不是 BoundedSet:调用方只拿它做一轮 has 查询就丢, * 没有长期持有,不需要上界。 */ 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 closed(B-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, // 淘汰计数持续增长说明上限设得太小 —— 那意味着会话上下文在被白白丢掉, // 而症状是「这条会话怎么突然不记得前面说过什么了」。 evictedSessions: sessionState.evicted, workers: [...running.values()].map((e) => ({ pid: e.child.pid, mailID: e.mailID, ageMs: Date.now() - e.startedAt, })), }); return { submit, routePermission, hasSession, forget, mailDrivenIDs, stop, stats, }; }