fix(pi): 等人点头的 worker 不再占并发额度,也不会被硬超时杀掉
实测(今天全 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,桥重启并重新心跳)。
This commit is contained in:
@ -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;
|
||||
}
|
||||
|
||||
@ -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();
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user