Files
MailUI4Agents/plugins/pi-mail-bridge/src/pool.mjs
JianFeeeee 453f451fbb fix(permission): 人类的备注必须到达模型 + 决策回执不再被当成新任务
用户报的「很严重的问题」:被拒绝的 agent 看不到授权备注,且看不到他发的回复邮件。
按数据查到了两个**真缺陷**,都在桥的权限回路上(不是猜测,三层证据)。

## 缺陷一:备注在桥内被连丢三处

网关其实一路都带着备注(`CreateDecisionMail(..., req.Note)` 把备注写进决策邮件正文,
SSE payload 里也有 `"note"`),但桥的三个环节只传 decision:
  index.mjs  `pool.routePermission(relayKey, String(data.decision))`
  pool.mjs   `child.send({type:'permission_decision', relayKey, decision})`
  worker.mjs `resolve(String(msg.decision))`
模型最终看到的只有 `用户拒绝了这次 bash 调用`(pi 会话转录逐字可查)。

现场:人类写「我说了让你拉取仓库到program下你听不懂吗」,模型不知道要改什么,
把同一条命令换个写法又问了 —— 会话里连问 **9 次**(22:16–22:26)。

## 缺陷二:决策回执照样被当"新任务"投递 + 等人的邮件被堵在后面

决策是**双通道**送达:SSE `permission_decision`(唤醒停放的 worker)+ 一封普通形状的
邮件("Re: 权限请求 - 拒绝")。以前两条都会起动作 ⇒ 同一件事被处理两次;而这条会话
的新邮件在 worker 停放期间只能排队。实测:人类 22:18:08 发出的更正
「不对,不是让你拉取到agentmail仓库,是让你拉取到program仓库!!」
直到 22:26:30(worker 回合结束)才被模型看到 —— **8 分钟**里它一直在错误的目录上打转。
转录里那封更正确实是模型自己 `read_mail` 读到的(不是没人给它)。

## 改动

- 网关:`CreateDecisionMail` 写 `mail_type='permission_decision'` —— 桥据此区分
  「控制面回执」与「新任务」。
- pi 桥(新增 `lib/denial-reason.js`、`lib/waiting-mails.js`):
  · 备注随决策一路透传到**模型看到的拒绝理由**(工具拦截与通知投递两条路都带);
  · 恢复停放的 worker 时,顺带把「等人期间新到、尚未标记已读」的邮件附进理由,
    模型当场就能改道(这正是那 8 分钟的洞);
  · 决策回执不再起新任务轮次(记进 deliveredMails);若决策事件尚未到达,
    退化为 B-4.3 的通知投递,且没有会话时不凭空新开。

## 判据

- `test/permission-note.test.mjs`:11 条(备注进理由、无备注不得凭空造说明、
  等人期间的邮件要点名 read_inbox、只挑本会话非权限类未交付的、上限、旧回包缺
  session_id 不能丢邮件、接线 8 处形状、判据自检)。
- **扰动验证**:把备注从 `pool.mjs` 的 send 里去掉 → 接线判据 2 条红;恢复 → 11 绿。
- 既有 pi 套件 420/420;server 10 包全绿(新增 1 条 Go 判据验决策邮件的类型与备注正文)。

## 现场证据(可复核)

- 桥日志:9 次 `权限 <key> 决策 同意/拒绝(决策人 jianf)已转交 worker`,全程不含备注;
  「worker 2135211 等待权限决策,让出并发额度(停放 1/5)」
- 会话转录:`{"toolName":"bash","content":[{"text":"用户拒绝了这次 bash 调用"}]}` ×6
2026-09-13 22:47:25 +08:00

405 lines
18 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* 投递工作进程池 —— 主进程侧的调度逻辑。
*
* # 它解决的问题
*
* 桥的主进程唯一的实时职责是读 SSE。模型工作放在主进程里跑会占满事件循环
* (pi 的会话装载是同步的:23MB 的会话文件 `SessionManager.open` 一次阻塞
* 118ms,实测;模型跑起来之后 SDK 内部还有大量同步工作),SSE 读循环停住,
* 后续邮件卡在 TCP 缓冲区,久到 Gateway 认为连接死了 → 重连 → 重放。
*
* 所以:**收到事件就派给一个子进程,主进程立刻回去读 SSE。**
*
* # 并发与串行的边界
*
* - **不同邮件会话并发**,上限 `maxWorkers`(默认 3)。上限的理由是内存
* (每个 worker 约 140MB RSS,实测)和对上游 provider 的并发请求数。
* - **等人点头的 worker 不占并发额度**(`maxParked` 单独限量)。见下。
* - **同一邮件会话串行**。这是正确性要求,不是限流:pi 没有任何锁机制,
* 它假定「一个文件一个持有者」。两个 worker 同时装载同一条会话文件,各自的
* 内存索引都看不见对方追加的行,算出的 parentId 指向对方不知道的 entry
* → 会话树分叉。串行还顺带保证了同一条线索里两封邮件的先后顺序。
*
* # 等人点头 ≠ 卡死:让出并发额度,也别被硬超时杀掉
*
* 实测:`maxWorkers=3` 会被**三个正在等人工授权**的 worker 吃满,于是新邮件只能
* 排队 —— 而人可能十分钟后才看邮箱。等待不是故障,不该占并发额度(额度管的是
* 内存与上游并发,等待中的 worker 两样都不占)。
*
* 同理,硬超时的用途是**回收卡死的进程**,等人点头不是卡死:等待期间暂停计时器、
* 决策到达后重启。否则 worker 会在人还在读邮件时被 SIGKILL,那次工具调用直接消失
* (人后来批了也没人接)。
*
* # 排队而不是拒绝
*
* 满载时邮件进 `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, maxParked = 5, 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();
}
/** 在干活的 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++) {
const job = queue[i];
// 同一会话已有 worker 在跑 → 跳过它,看后面有没有别的会话可以先跑。
// 不能 break:那会让一条慢会话把所有别的会话都堵住(正是要修的病)。
if (running.has(job.key)) continue;
if (activeCount() >= 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'],
});
const entry = {
child, mailID: job.data?.mail_id || '', key: job.key,
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) => {
// 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);
// 没收到 `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);
// 让出并发额度:等人工决策的 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 一封一进程,不存的话下一封
// 邮件又问一遍,那个选项就是在骗人。
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, note = '', freshMails = []) {
const key = permissionRoutes.get(relayKey);
if (!key) return false;
const entry = running.get(key);
if (!entry) {
permissionRoutes.delete(relayKey);
return false;
}
permissionRoutes.delete(relayKey);
if (entry.parked) {
// 决策到了 → 它重新开始干活:回到并发额度,并给一个完整超时窗口
entry.parked = false;
armTimeout(entry);
log(`worker ${entry.child.pid} 收到决策,恢复占用并发额度(停放 ${parkedCount()}/${maxParked})`);
}
// 备注与「等人期间新到的邮件」必须一起送到:以前只传 decision,于是人类写
// 「我说了让你拉取仓库到program下你听不懂吗」,模型只看到「拒绝」,
// 转头把同一条命令又问了一遍(2026-09-13 实测连问 9 次)。
entry.child.send({ type: 'permission_decision', relayKey, decision, note, freshMails });
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 };
}