From 49ec22f91622230305174d188f884d3399a28505 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 13 Sep 2026 13:14:56 +0800 Subject: [PATCH] =?UTF-8?q?fix(pi):=20=E7=AD=89=E4=BA=BA=E7=82=B9=E5=A4=B4?= =?UTF-8?q?=E7=9A=84=20worker=20=E4=B8=8D=E5=86=8D=E5=8D=A0=E5=B9=B6?= =?UTF-8?q?=E5=8F=91=E9=A2=9D=E5=BA=A6=EF=BC=8C=E4=B9=9F=E4=B8=8D=E4=BC=9A?= =?UTF-8?q?=E8=A2=AB=E7=A1=AC=E8=B6=85=E6=97=B6=E6=9D=80=E6=8E=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 实测(今天全 Agent 演练时撞到的):pi 的 `maxWorkers=3` 被**三个正在等人工授权**的 worker 吃满,于是新邮件只能排队 —— 而人可能十分钟后才看邮箱。清掉卡住的请求后队列 立即排空,机制本身没错,错的是"等待"被当成了"在干活"。 两处改动(`src/pool.mjs`): 1. **等待期间让出并发额度**。worker 发 `permission_pending` 时把它的 entry 标为 parked,`pump()` 只数"在干活"的(`activeCount()`)。停放另有上限 `maxParked` (默认 5,防止内存无界:每 worker 约 140MB);超出后仍占额度并打日志说明。 决策到达(`routePermission`)时解除停放,回到额度里。 2. **等待期间暂停硬超时**。硬超时的用途是回收**卡死**的进程,而等人点头不是卡死: 停放时清掉计时器,决策到达后重新起一个完整窗口。否则 worker 会在人还在读邮件时 被 SIGKILL —— 那次工具调用直接消失,人后来批了也没人接(这类"批了没反应"的现象 与此吻合)。 判据(`test/pool.test.mjs` 新增 3 条): · ★ 等待授权的 worker 让出额度:另一个会话的邮件必须能开跑 · ★ 等人点头期间(800ms > 300ms 硬超时)不得被强杀,且恢复后能跑完 · 停放有上限:超出后仍占额度(不无限超发) **扰动验证**:整块回退到 HEAD → 3 条全红;修复后 3/3。全量 pi 套件 420/420。 已部署(快照 20260913-131216,桥重启并重新心跳)。 --- plugins/pi-mail-bridge/src/pool.mjs | 74 +++++++++++++++++++---- plugins/pi-mail-bridge/test/pool.test.mjs | 61 +++++++++++++++++++ 2 files changed, 122 insertions(+), 13 deletions(-) diff --git a/plugins/pi-mail-bridge/src/pool.mjs b/plugins/pi-mail-bridge/src/pool.mjs index 015703f..041b327 100644 --- a/plugins/pi-mail-bridge/src/pool.mjs +++ b/plugins/pi-mail-bridge/src/pool.mjs @@ -14,11 +14,22 @@ * * - **不同邮件会话并发**,上限 `maxWorkers`(默认 3)。上限的理由是内存 * (每个 worker 约 140MB RSS,实测)和对上游 provider 的并发请求数。 + * - **等人点头的 worker 不占并发额度**(`maxParked` 单独限量)。见下。 * - **同一邮件会话串行**。这是正确性要求,不是限流:pi 没有任何锁机制, * 它假定「一个文件一个持有者」。两个 worker 同时装载同一条会话文件,各自的 * 内存索引都看不见对方追加的行,算出的 parentId 指向对方不知道的 entry * → 会话树分叉。串行还顺带保证了同一条线索里两封邮件的先后顺序。 * + * # 等人点头 ≠ 卡死:让出并发额度,也别被硬超时杀掉 + * + * 实测:`maxWorkers=3` 会被**三个正在等人工授权**的 worker 吃满,于是新邮件只能 + * 排队 —— 而人可能十分钟后才看邮箱。等待不是故障,不该占并发额度(额度管的是 + * 内存与上游并发,等待中的 worker 两样都不占)。 + * + * 同理,硬超时的用途是**回收卡死的进程**,等人点头不是卡死:等待期间暂停计时器、 + * 决策到达后重启。否则 worker 会在人还在读邮件时被 SIGKILL,那次工具调用直接消失 + * (人后来批了也没人接)。 + * * # 排队而不是拒绝 * * 满载时邮件进 `queue`,有 worker 空出来就派。丢掉邮件是不可接受的: @@ -78,7 +89,7 @@ const WORKER_PATH = fileURLToPath(new URL('./worker.mjs', import.meta.url)); */ export function createWorkerPool({ log, config, onReconfigure, - maxWorkers = 3, workerMaxMs = 600_000, maxAttempts = 3, workerPath = WORKER_PATH, + maxWorkers = 3, maxParked = 5, workerMaxMs = 600_000, maxAttempts = 3, workerPath = WORKER_PATH, }) { /** 正在跑的 worker:mailSessionKey -> {child, mailID, startedAt, timer} */ const running = new Map(); @@ -117,6 +128,29 @@ export function createWorkerPool({ pump(); } + /** 在干活的 worker 数(等人的不算 —— 见文件头「等人点头 ≠ 卡死」)。 */ + function activeCount() { + let n = 0; + for (const e of running.values()) if (!e.parked) n += 1; + return n; + } + + const parkedCount = () => { + let n = 0; + for (const e of running.values()) if (e.parked) n += 1; + return n; + }; + + /** 重启硬超时(停放或被派活时都要有一次完整窗口)。 */ + function armTimeout(entry) { + if (entry.timer) clearTimeout(entry.timer); + entry.timer = setTimeout(() => { + log(`worker ${entry.child.pid} 处理 ${entry.mailID} 超过 ${workerMaxMs / 1000}s,强杀`); + try { entry.child.kill('SIGKILL'); } catch { /* 已经死了 */ } + }, workerMaxMs); + if (typeof entry.timer.unref === 'function') entry.timer.unref(); + } + function pump() { if (stopped) return; for (let i = 0; i < queue.length; i++) { @@ -124,7 +158,7 @@ export function createWorkerPool({ // 同一会话已有 worker 在跑 → 跳过它,看后面有没有别的会话可以先跑。 // 不能 break:那会让一条慢会话把所有别的会话都堵住(正是要修的病)。 if (running.has(job.key)) continue; - if (running.size >= maxWorkers) return; + if (activeCount() >= maxWorkers) return; queue.splice(i, 1); i--; spawn(job); @@ -139,24 +173,20 @@ export function createWorkerPool({ 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, + startedAt: Date.now(), timer: null, settled: false, parked: false, }; running.set(job.key, entry); + // 硬超时:worker 卡死(模型不返回、权限等不到决策而主进程也没收到事件) + // 时必须能回收,否则那条会话的后续邮件永远排队。 + armTimeout(entry); child.on('message', (msg) => onWorkerMessage(entry, msg)); child.on('exit', (code, signal) => { - clearTimeout(timer); + // entry.timer 可能为 null(停放期间计时器是停的)——clearTimeout(null) 无副作用 + clearTimeout(entry.timer); running.delete(job.key); for (const [rk, k] of permissionRoutes) if (k === job.key) permissionRoutes.delete(rk); @@ -232,9 +262,21 @@ export function createWorkerPool({ state.cwd = msg.cwd; sessionState.set(entry.key, state); return; - case 'permission_pending': + case 'permission_pending': { permissionRoutes.set(msg.relayKey, entry.key); + // 让出并发额度:等人工决策的 worker 不是"在干活",不该把池子占满 + // (实测:三个等待授权的 worker 会把 maxWorkers=3 用光,新邮件只能排队)。 + if (!entry.parked && parkedCount() < maxParked) { + entry.parked = true; + if (entry.timer) clearTimeout(entry.timer); + entry.timer = null; // 等人的时间不计入硬超时(见文件头) + log(`worker ${entry.child.pid} 等待权限决策,让出并发额度(停放 ${parkedCount()}/${maxParked})`); + pump(); + } else if (!entry.parked) { + log(`worker ${entry.child.pid} 等待权限决策,但停放额度已满(${maxParked}),仍占并发额度`); + } return; + } case 'permission_grant': // 「一直同意」必须跨 worker 活着:worker 一封一进程,不存的话下一封 // 邮件又问一遍,那个选项就是在骗人。 @@ -275,6 +317,12 @@ export function createWorkerPool({ return false; } permissionRoutes.delete(relayKey); + if (entry.parked) { + // 决策到了 → 它重新开始干活:回到并发额度,并给一个完整超时窗口 + entry.parked = false; + armTimeout(entry); + log(`worker ${entry.child.pid} 收到决策,恢复占用并发额度(停放 ${parkedCount()}/${maxParked})`); + } entry.child.send({ type: 'permission_decision', relayKey, decision }); return true; } diff --git a/plugins/pi-mail-bridge/test/pool.test.mjs b/plugins/pi-mail-bridge/test/pool.test.mjs index 932be4c..1672469 100644 --- a/plugins/pi-mail-bridge/test/pool.test.mjs +++ b/plugins/pi-mail-bridge/test/pool.test.mjs @@ -92,6 +92,7 @@ function makePool(opts = {}) { config: () => ({ turnTimeoutMs: 1000, ...(opts.config || {}) }), onReconfigure: opts.onReconfigure || (() => {}), maxWorkers: opts.maxWorkers ?? 2, + maxParked: opts.maxParked ?? 5, workerMaxMs: opts.workerMaxMs ?? 5000, maxAttempts: opts.maxAttempts, workerPath: opts.workerPath || STUB, @@ -418,3 +419,63 @@ test('重投上限之下不会无限重投(maxAttempts=1 就是不重投)', assert.equal(jobs(lines).length, 1, 'maxAttempts=1 时只跑一次'); }); + +// ─── 等人点头 ≠ 卡死(2026-09-13 实测后加的)──────────────────────── +// +// 实测故障:maxWorkers=3 被三个"正在等人工授权"的 worker 吃满,于是新邮件只能 +// 排队 —— 而人可能十分钟后才看邮箱。等待不是故障:它既不占内存活动量也不打上游, +// 不该占并发额度;同理它也不是"卡死",不该被硬超时回收(否则人还在读邮件, +// 那次工具调用就被 SIGKILL 了,人后来批了也没人接)。 + +test('★ 等人工授权的 worker 让出并发额度:别的会话不再被它堵住', async () => { + const { pool, lines } = makePool({ maxWorkers: 1 }); + pool.submit('mail', { mail_id: 'waiting', session_id: 'WAIT', __pending: 'rk-wait' }); + await until(() => pool.stats().running === 1, 1500); + + // 另一个会话的邮件:修复前额度被等待者占着,它永远不会开跑 + pool.submit('mail', { mail_id: 'other', session_id: 'OTHER', __hold: 30 }); + const started = await until(() => jobs(lines).some((j) => j.mailID === 'other'), 2000); + assert.ok(started, '等待授权的 worker 不该占着额度 —— 另一个会话必须能开跑'); + + // 决策到达 → 等待者恢复、继续跑完 + assert.equal(pool.routePermission('rk-wait', '同意'), true); + const drained = await until(() => pool.stats().running === 0, 4000); + pool.stop(); + assert.ok(drained, '收到决策后应正常结束'); + assert.ok(lines.some((l) => l.includes('DECISION rk-wait=同意')), `worker 应收到决策,实际:\n${lines.join('\n')}`); +}); + +test('★ 等人点头期间不被硬超时杀掉(等待不是卡死)', async () => { + const { pool, lines } = makePool({ maxWorkers: 1, workerMaxMs: 300 }); + pool.submit('mail', { mail_id: 'hold-perm', session_id: 'HP', __pending: 'rk-hp' }); + await until(() => pool.stats().running === 1, 1000); + + await sleep(800); // 远超 300ms 硬超时 + assert.equal(pool.stats().running, 1, '等待决策的 worker 不该被硬超时回收'); + assert.ok(!lines.join('\n').includes('强杀'), '等待期间的强杀不该发生'); + + pool.routePermission('rk-hp', '同意'); + const drained = await until(() => pool.stats().running === 0, 3000); + pool.stop(); + assert.ok(drained, '恢复后应能正常跑完(超时窗口重启)'); +}); + +test('停放有上限:超出后仍占用额度(不无限超发)', async () => { + const { pool, lines } = makePool({ maxWorkers: 3, maxParked: 1 }); + try { + pool.submit('mail', { mail_id: 'p1', session_id: 'P1', __pending: 'rk-p1' }); + pool.submit('mail', { mail_id: 'p2', session_id: 'P2', __pending: 'rk-p2' }); + await until(() => jobs(lines).length >= 2, 2500); + await sleep(250); + + assert.ok( + lines.some((l) => l.includes('停放额度已满')), + `第二个等待者应记录"额度已满"并继续占额度,实际:\n${lines.join('\n')}` + ); + for (const rk of ['rk-p1', 'rk-p2']) pool.routePermission(rk, '同意'); + await until(() => pool.stats().running === 0, 3000); + } finally { + // 断言失败时也要收掉子进程:否则挂住的 worker 会拖住整个测试文件(实测过) + pool.stop(); + } +});