# 发生了什么 pi-lens 内置「安全格式化」:它会自动安装 biome 并对**编辑过的文件**跑 `biome format --write`。本机原先没有任何 biome 配置,于是 biome 用它自己的 默认值 —— tab 缩进 + 双引号 —— 把文件整体重写。 我在19a3161那次提交里用了 `git add -A`,把这批与功能无关的重排一起扫了进去: 约 7000 行改动散落在 20 个文件上,使那次提交无法审查,还掩盖了 server/internal/handler/permission.go 的一处删行(实为文件末尾空行,无代码丢失)。 # 为什么是「关掉」而不是「配置成我们的风格」 试过把缩进/引号/lineWidth 全部对齐本仓库习惯(biome.json + space/2/single/ lineWidth 120):`biome format --write` 仍然改动 17 个文件。原因是本仓库从未按 biome 的规则排版过 —— 注释按语义换行、数组与调用按可读性手工折行, 这些无法由格式化器还原。也就是说只要格式化器开着,每次编辑都会产生与内容无关的 大面积 diff,把真正的改动埋掉。 因此 biome.jsonc 里 formatter 与 linter 都关闭:本仓库的静态检查由 tsc / go vet / tree-sitter / ast-grep 与各自测试套件承担,不引入会改动无关行的 自动修复。 (pi-lens 这一版把 format 服务的 enabled 硬编码为 true,没有配置开关, 所以只能在仓库侧用 biome 配置让它不动文件;已验证 `biome format --write` 对这些文件零改动。) # 本提交内容 把19a3161里除「有意改动」外的 20 个文件还原到重排前的样子。19a3161中真正有意的改动是 deploy/install.sh 的扩展注册与 plugins/pi-mail-bridge/extension/index.ts 新文件,两者原样保留。 验证:Go 全量、三桥插件(320/362/409)、前端 196 全绿; `biome format --write` 对还原后的文件零改动。
354 lines
16 KiB
JavaScript
354 lines
16 KiB
JavaScript
/**
|
||
* 投递工作进程池 —— 主进程侧的调度逻辑。
|
||
*
|
||
* # 它解决的问题
|
||
*
|
||
* 桥的主进程唯一的实时职责是读 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 };
|
||
}
|