- worker.mjs: mailContext 加 permissionMode 字段(从 SSE payload 读入) - tool_call hook 增加档位判定: full → 不拦截任何工具(直接 return) plan → 被守卫工具(bash/write/edit)一律 block + 返回原因说明 workspace → 走原有问人流程(不变) - index.mjs 心跳上报 mode_enforcement: 'native' - 导入 normalizeMode/MODE_FULL/MODE_PLAN 从 lib/permission-mode.js - 377 测试全过
622 lines
28 KiB
JavaScript
622 lines
28 KiB
JavaScript
#!/usr/bin/env node
|
||
/**
|
||
* pi-mail-bridge 的**投递工作进程** —— 一封邮件一个,跑完就退。
|
||
*
|
||
* # 为什么必须是独立进程
|
||
*
|
||
* 桥的主进程要一直读 SSE。而 pi 的会话装载是**同步**的:
|
||
* `SessionManager.open()` 走 `openSync` + `readSync` 循环把整个 `.jsonl`
|
||
* 读进内存并逐行 JSON.parse(SDK core/session-manager.js 的
|
||
* loadEntriesFromFile)。实测本机最大那条会话 23MB,`open` 一次
|
||
* **阻塞事件循环 118ms**;模型跑起来之后 SDK 内部还有大量同步工作。
|
||
* 全都发生在主线程上,SSE 读循环在那期间完全停住 —— 后续邮件卡在 TCP
|
||
* 缓冲区里,久到 Gateway 认为连接死了,重连又触发重放。
|
||
*
|
||
* 把模型工作搬进子进程后,主进程只剩「收事件 → 去重 → 分派」,
|
||
* 实测同样的活在 fork 出的子进程里跑,主进程事件循环阻塞 0ms。
|
||
*
|
||
* # 为什么是 child_process 而不是 worker_threads
|
||
*
|
||
* 两者实测都能建起 AgentSession。选进程的理由是**隔离**:
|
||
* 模型会跑 bash/write/edit,一次 OOM 或 uncaughtException 不该带走
|
||
* 整座桥;而 `worker_threads` 与主线程共享堆和进程生命周期。
|
||
* 代价是每个 worker 约 140MB RSS 和 ~500ms 启动,由主进程的并发上限约束。
|
||
*
|
||
* # 一封邮件一个进程带来的简化
|
||
*
|
||
* 原来的「接管会话短暂持有」(adopted / adoptTimers / releaseAdopted)
|
||
* 整套机制没有了:worker 退出**就是**释放,而且对普通会话和接管会话
|
||
* 一视同仁 —— 每封邮件都是「open → 跑一轮 → 还回去」。
|
||
* pi 假定「一个文件一个持有者」,主进程按会话串行分派保证了这一点。
|
||
*
|
||
* # 进程边界上传什么
|
||
*
|
||
* 主进程持有的是**路径与标量**(sessionFile / cwd / piSessionId / 授权表),
|
||
* 不是 AgentSession 对象 —— 那东西跨不了进程。worker 每次从 sessionFile
|
||
* 重新装载,拿到的是包含 TUI 期间写入的全部历史。
|
||
*
|
||
* 协议见 src/pool.mjs 顶部。
|
||
*/
|
||
|
||
import { existsSync } from 'node:fs';
|
||
import { homedir } from 'node:os';
|
||
import { join } from 'node:path';
|
||
import { ModelRuntime } from '@earendil-works/pi-coding-agent';
|
||
|
||
import { GatewayClient } from './gateway.mjs';
|
||
import { createMailTools } from './tools.mjs';
|
||
import { openSession, runTurn } from './session-pool.mjs';
|
||
import { buildMailPrompt, lastAssistantText, replySubject, relayKeyFor, describeError } from './turn.mjs';
|
||
import { planNamingSync, planWriteBack } from './naming.mjs';
|
||
import { resolveWorkspaceCwd, ensureCwd } from '../lib/workspace.js';
|
||
import { modelAttemptOrder, renderFailureReport } from '../lib/model-scope.js';
|
||
import { explicitSends, shouldSkipAutoRelay } from '../lib/relay-dedup.js';
|
||
import { autoRelayDecision } from '../lib/relay-policy.js';
|
||
import { adoptedSessionID, adoptMissingMessage } from '../lib/adopt.js';
|
||
import { isApproval, isAlwaysDecision } from '../lib/permission-grants.js';
|
||
import { normalizeMode, MODE_FULL, MODE_PLAN } from '../lib/permission-mode.js';
|
||
import { clampRelayKey, isPermanentFailure } from '../lib/relay-key.js';
|
||
|
||
// ─── 与主进程的通道 ───
|
||
|
||
/** 日志一律回传主进程:worker 的 stdout 会混在一起,加 pid 前缀才分得清。 */
|
||
const log = (...args) => send({ type: 'log', line: args.join(' ') });
|
||
|
||
function send(msg) {
|
||
try {
|
||
process.send?.(msg);
|
||
} catch {
|
||
/* 主进程已经走了,这条日志没有去处 */
|
||
}
|
||
}
|
||
|
||
/** relay_key -> resolve;决策由主进程的 SSE 收到后路由进来。 */
|
||
const pending = new Map();
|
||
|
||
/**
|
||
* 这一轮里人点过「一直同意」的工具。
|
||
*
|
||
* 初值由主进程在 job 里给(跨 worker 持久),新增的回报给主进程。
|
||
* 不用 lib/permission-grants.js 的 store:那份按 (会话, 工具) 存,
|
||
* 而 worker 只服务一条会话,一个 Set 就够,且要能整体回传。
|
||
*/
|
||
const grants = new Set();
|
||
|
||
// ─── 单封邮件的全部状态 ───
|
||
//
|
||
// worker 只处理一封邮件、只碰一条会话,所以这些原来在主进程里
|
||
// 按 sessionId 分桶的 Map 在这里都退化成单个变量。
|
||
|
||
let job = null;
|
||
let client = null;
|
||
let modelRuntime = null;
|
||
let piSessionId = '';
|
||
let mailContext = { replyTo: '', subject: '', mailID: '', permissionMode: 'workspace' };
|
||
let lastSyncedName = '';
|
||
let relayedKey = '';
|
||
let finished = false;
|
||
|
||
// ─── 权限钩子(B-8)───
|
||
|
||
/**
|
||
* 与主进程版本逐条对应,差别只有两处:
|
||
* - 「是不是邮件驱动的会话」不必查表:worker 只为邮件而存在。
|
||
* - 等决策的 promise 由主进程通过 IPC 唤醒,而不是本进程的 SSE。
|
||
*
|
||
* 决策等待期间**只有这个 worker 停住**,主进程照常读 SSE、照常给别的会话
|
||
* 派活 —— 这正是原来最难受的一处:权限询问会让整座桥不再收信。
|
||
*/
|
||
function permissionExtension() {
|
||
const GUARDED = new Set(['bash', 'write', 'edit']);
|
||
|
||
return (pi) => {
|
||
pi.on('tool_call', async (event, ctx) => {
|
||
// ── 档位判定 ──
|
||
// plan: 被守卫的工具一律直接拒绝(该档语义是「不动手」,没什么可问人的)
|
||
// full: 不拦截(已声明全权)
|
||
// workspace: 走原有问人流程
|
||
const mode = normalizeMode(mailContext.permissionMode);
|
||
if (mode === MODE_FULL) return; // full 档不拦任何工具
|
||
if (mode === MODE_PLAN && GUARDED.has(event.toolName)) {
|
||
return {
|
||
block: true,
|
||
reason: `plan 档下不允许执行 ${event.toolName}。本档只允许读与查,请把方案写在回信里。如需动手请让发件人把档位改成 workspace。`,
|
||
};
|
||
}
|
||
if (!GUARDED.has(event.toolName)) return;
|
||
|
||
const sid = ctx?.sessionManager?.getSessionId?.() || '';
|
||
// 只管自己那条会话。worker 里不该出现第二条,出现了说明有 bug ——
|
||
// 让位(返回 undefined)比拦错一个安全。
|
||
if (piSessionId && sid && sid !== piSessionId) return;
|
||
|
||
// 人点过「一直同意」→ 直接放行。必须在 POST 之前:否则每条命令都发一封
|
||
// 邮件,那个选项形同虚设(生产实测同一条会话被问了 15 次 bash)。
|
||
if (grants.has(event.toolName)) return;
|
||
|
||
// relay_key 用 pi 给的 toolCallId(B-8.1):服务端随决策事件回传它。
|
||
//
|
||
// clampRelayKey 不是防御性冗余:启用 extended thinking 时 Bedrock 把
|
||
// 思考签名拼进 toolCallId,实测长到 437 ~ 13601 字节,键直接超服务端
|
||
// 160 字节列宽 → 400。同一条会话里长短两种形态混着出现,于是权限询问
|
||
// 随机成功随机失败(生产日志 21:54:02 失败、21:54:24 同会话成功)。
|
||
const relayKey = clampRelayKey(`${sid || piSessionId}:${event.toolCallId}`);
|
||
|
||
try {
|
||
// 不传 `to`:决策人由服务端按 会话 owner → 线索里最近的人类 → 409
|
||
// 解析。插件只有本地上下文,猜不出「这条 Agent 链最初是谁派的活」。
|
||
await client.post('/permission/request', {
|
||
question: `是否允许执行 ${event.toolName}?`,
|
||
options: ['同意', '一直同意', '拒绝'],
|
||
context: [
|
||
describeToolCall(event),
|
||
mailContext.subject ? `\n触发任务:${mailContext.subject}` : '',
|
||
mailContext.replyTo ? `任务来自:${mailContext.replyTo}` : '',
|
||
].filter(Boolean).join('\n'),
|
||
session_id: job.data?.session_id || '',
|
||
relay_key: relayKey,
|
||
});
|
||
} catch (e) {
|
||
// 409 = 服务端判定这条任务链上没有人类,永远不会有人来点头。
|
||
// 当场 block 并把服务端建议原文当 reason:模型从工具报错里看到
|
||
// 「没人可问,换不需要权限的方式」才能自己改道,挂死时连重试机会都没有。
|
||
if (e?.status === 409) {
|
||
const b = e.body || {};
|
||
const reason = [
|
||
b.error || '权限询问无法送达:这条任务链上没有人类用户',
|
||
b.detail || '',
|
||
b.suggestion || '',
|
||
].filter(Boolean).join('\n');
|
||
log(`权限询问无人可投,当场拒绝 ${relayKey}:${b.error || ''}`);
|
||
return { block: true, reason };
|
||
}
|
||
|
||
// 其余 4xx(400 / 401 / 403 / 404 / 422…)同样永远不会因重试成功。
|
||
//
|
||
// 这里原来一律「让位给本地决策」,而邮件驱动的 worker 没有 TUI ——
|
||
// 让位等于守卫消失。生产实测:relay_key 过长报 400 被当暂时失败,
|
||
// 那次 bash 在无人批准的情况下执行了(21:54:02 让位,同会话
|
||
// 21:54:24 另一次 key 正常,于是权限询问被随机吞掉)。
|
||
//
|
||
// fail closed:宁可让模型看到「权限系统坏了」并自己改道,
|
||
// 也不能悄悄放行一条没人看过的命令。
|
||
if (isPermanentFailure(e)) {
|
||
const detail = describeError(e);
|
||
log(`权限转发遇到永久失败(HTTP ${e?.status}),当场拒绝:${detail}`);
|
||
return {
|
||
block: true,
|
||
reason: [
|
||
`无法把 ${event.toolName} 的授权请求送达给人类:${detail}`,
|
||
'这是一个不会因重试而改变的失败(请求本身被服务端拒绝)。',
|
||
'请改用不需要授权的方式完成,或在回信里说明需要人工执行哪一步。',
|
||
].join('\n'),
|
||
};
|
||
}
|
||
|
||
// 暂时失败(5xx / 408 / 429 / 网络层)→ 让位给 pi 本地决策(B-8.2)。
|
||
log(`权限转发暂时失败,让位给本地决策: ${describeError(e)}`);
|
||
return;
|
||
}
|
||
|
||
log(`权限询问已发出(${event.toolName},key=${relayKey}),等待决策…`);
|
||
send({ type: 'permission_pending', relayKey });
|
||
const decision = await new Promise((resolve) => pending.set(relayKey, resolve));
|
||
|
||
// fail closed(B-9.2 / N-9):只有明确同意才放行。
|
||
if (isApproval(decision)) {
|
||
// 「一直同意」要真的记住,否则这个选项在骗人。判定交给 isAlwaysDecision ——
|
||
// 「同意」是单次授权,把它当 always 会放行人没看过的后续命令。
|
||
if (isAlwaysDecision(decision)) {
|
||
grants.add(event.toolName);
|
||
send({ type: 'permission_grant', toolName: event.toolName });
|
||
log(`本会话的 ${event.toolName} 已获「一直同意」,后续不再询问`);
|
||
}
|
||
log(`权限 ${relayKey} 获批(${decision}),放行 ${event.toolName}`);
|
||
return;
|
||
}
|
||
return {
|
||
block: true,
|
||
reason: `用户${decision === 'shutdown' ? '未及决策(桥已关停)' : `拒绝了这次 ${event.toolName} 调用`}`,
|
||
};
|
||
});
|
||
};
|
||
}
|
||
|
||
/** 把一次工具调用摘要成人能判断的文本(B-8.4)。 */
|
||
function describeToolCall(event) {
|
||
const input = event?.input ?? {};
|
||
if (event.toolName === 'bash') return `命令:\n${String(input.command ?? '').slice(0, 800)}`;
|
||
if (event.toolName === 'write' || event.toolName === 'edit') {
|
||
return `文件:${input.file_path ?? input.path ?? '(未给出)'}`;
|
||
}
|
||
return JSON.stringify(input).slice(0, 800);
|
||
}
|
||
|
||
// ─── 会话装载 ───
|
||
|
||
/**
|
||
* 没有可用 `to_workspace` 时的兜底目录。
|
||
*
|
||
* 与 DSH 的 `mailSessionFallback` 同构,但目录名是 `.pi`:那个函数在 lib/ 下
|
||
* (三平台逐字节相同),写死了 `.dsh`,不能为 pi 改 —— pi 的会话落进 `~/.dsh/`
|
||
* 会让人以为是 DSH 在干活。
|
||
*/
|
||
function piMailFallback(sessionKey) {
|
||
return join(homedir(), '.pi', 'mail-sessions', String(sessionKey || 'default'));
|
||
}
|
||
|
||
/**
|
||
* 找到这封邮件该落进的会话文件,装载它。
|
||
*
|
||
* 三条路,优先级从高到低:
|
||
* 1. 主进程给了 sessionFile —— 这条邮件会话之前有 worker 跑过,接着谈。
|
||
* 2. 服务端说接管了平台会话(人在 TUI 里开的那条)→ 按 platform id 找文件。
|
||
* 3. 都没有 → 新建,cwd 取地址的 path 位(B-3.1)。
|
||
*
|
||
* 返回 `reused` 供提示词与模型降级判断:有历史的会话不重新自我介绍,
|
||
* 也不做模型降级(换模型要换会话,会丢掉整条上下文,而上下文正是
|
||
* 发件人指定这条会话的原因)。
|
||
*/
|
||
async function loadSession(mailTools) {
|
||
const data = job.data;
|
||
const given = job.session?.sessionFile;
|
||
|
||
if (given && existsSync(given)) {
|
||
const cwd = job.session.cwd || resolveWorkspaceCwd(
|
||
data.to_workspace, piMailFallback(data.session_id)).cwd;
|
||
const opened = await openSession({
|
||
cwd, modelRuntime, customTools: mailTools,
|
||
extension: permissionExtension(), sessionFile: given,
|
||
});
|
||
return { ...opened, cwd, reused: true };
|
||
}
|
||
|
||
const adoptID = adoptedSessionID(data);
|
||
if (adoptID) {
|
||
// 用 sessionScanner 而不是 `SessionManager.listAll()`:这里只要 id → path,
|
||
// 而那两个字段全在会话文件的**首行** header 里。listAll 为了拿它们会把
|
||
// 每个 .jsonl 的每一行读进来并 JSON.parse,还把所有消息正文拼成一个大字符串
|
||
// (本机 115 个文件 / 145MB 实测:1431ms、堆里瞬时 240MB)。
|
||
//
|
||
// worker 是短命进程,拿不到跨拍缓存的好处,但冷启动也依然便宜得多:
|
||
// 没有任何一行巨大的 message 被 materialize(实测单行最长 2.63MB)。
|
||
const { createSessionScanner } = await import('./session-scan.mjs');
|
||
const { getAgentDir } = await import('@earendil-works/pi-coding-agent');
|
||
const all = await createSessionScanner({
|
||
sessionsDir: join(getAgentDir(), 'sessions'),
|
||
}).scan();
|
||
const info = all.find((e) => e?.id === adoptID);
|
||
if (!info?.path) {
|
||
// 镜像是快照,可以过期。**不能**退回「新建一条」—— 那会让人在 TUI 里
|
||
// 看不到这封邮件带来的对话,而那正是接管的目的(N-8:静默改语义比报错糟)。
|
||
throw new Error(adoptMissingMessage(adoptID, '磁盘上已无这个会话文件'));
|
||
}
|
||
// cwd 取会话自己的(SessionInfo.cwd 来自持久化 header);
|
||
// 老会话的 cwd 是空串,那种情况退回地址里的 path 位。
|
||
const { cwd } = resolveWorkspaceCwd(
|
||
info.cwd || data.to_workspace, piMailFallback(data.session_id));
|
||
const opened = await openSession({
|
||
cwd, modelRuntime, customTools: mailTools,
|
||
extension: permissionExtension(), sessionFile: info.path,
|
||
});
|
||
log(`接管 pi 会话 ${adoptID}(cwd=${cwd},文件 ${info.path})`);
|
||
return { ...opened, cwd, reused: true };
|
||
}
|
||
|
||
// 目录不存在时**不创建**(N-2:笔误会在磁盘上落下真目录,而 Agent 在
|
||
// 里面一无所获),拒绝相对路径(N-3)。
|
||
const { cwd, grouped } = resolveWorkspaceCwd(
|
||
data.to_workspace, piMailFallback(data.session_id));
|
||
if (!grouped && data.to_workspace) {
|
||
log(`工作目录 ${data.to_workspace} 不可用,回退到 ${cwd}`);
|
||
}
|
||
ensureCwd(cwd, grouped);
|
||
const opened = await openSession({
|
||
cwd, modelRuntime, customTools: mailTools, extension: permissionExtension(),
|
||
});
|
||
log(`新建 pi 会话 ${opened.session.sessionId}(cwd=${cwd})`);
|
||
return { ...opened, cwd, reused: false };
|
||
}
|
||
|
||
// ─── 命名一致(C-11 / W-7)───
|
||
|
||
async function syncNaming(session, platformName) {
|
||
const mailSessionID = job.data?.session_id;
|
||
if (!mailSessionID) return;
|
||
|
||
const plan = planNamingSync({
|
||
platformName,
|
||
mailSubject: mailContext.subject,
|
||
lastSynced: lastSyncedName,
|
||
});
|
||
if (plan.skip) return;
|
||
|
||
// 先记指纹再发请求:响应回来时 setSessionName 会再次触发
|
||
// session_info_changed,这一步是防自激循环的关键。
|
||
lastSyncedName = plan.signature;
|
||
send({ type: 'name_synced', signature: plan.signature });
|
||
|
||
const res = await client.post(`/sessions/${mailSessionID}/sync`, {
|
||
alias: plan.alias, title: plan.title,
|
||
});
|
||
|
||
const back = planWriteBack({ finalAlias: res?.alias, currentPiName: session.sessionName });
|
||
log(`命名同步 ${piSessionId}: alias=${res?.alias || '(未变)'} 来源=${plan.source}`);
|
||
|
||
if (back.write) {
|
||
// 顺序要紧:先更新指纹,再改名。setSessionName **同步**触发
|
||
// session_info_changed(实测),本函数会被重入;指纹后更新的话重入那次
|
||
// 看到旧指纹,于是又打一次 sync —— 每条会话两次内容相同的请求。
|
||
lastSyncedName = `platform:${back.name}|${back.name}`;
|
||
send({ type: 'name_synced', signature: lastSyncedName });
|
||
// 只用 setSessionName(走 pi 自己的写入路径)。绝不自己拼路径写会话文件:
|
||
// 首条 assistant 消息落盘前文件还不存在,pi 首次落盘用 openSync(file,"wx"),
|
||
// 抢先创建会让它抛 EEXIST(实测)。
|
||
session.setSessionName(back.name);
|
||
log(`别名回写 pi:${back.name}(${back.reason})`);
|
||
}
|
||
}
|
||
|
||
// ─── 自动转发(B-5)───
|
||
|
||
/**
|
||
* 把这一轮的结论转回发件人。
|
||
*
|
||
* 与主进程版本的差别:不再挂在 `agent_end` 订阅上,而是在 `runTurn` 返回后
|
||
* **确定性地**调一次。原来必须靠事件是因为主进程的 60 秒超时会先返回、
|
||
* 会话还在跑;worker 没有那个约束(它就为这封邮件活着),等真结束再转发。
|
||
*
|
||
* @param {boolean} adopted 接管会话跳过命名同步 —— 它的别名是人从补全里
|
||
* 选中的平台 slug,同步会双向改坏:撞名时 Gateway 加后缀,定稿别名又回写进
|
||
* pi 会话文件,于是下次心跳上报的 slug 变成带后缀那个。实测撞出来过一次。
|
||
*/
|
||
async function relaySummary(session, sessionManager, adopted) {
|
||
if (!adopted) {
|
||
await syncNaming(session, session.sessionName)
|
||
.catch((e) => log(`命名同步失败: ${describeError(e)}`));
|
||
}
|
||
|
||
// **只给人类来信自动转发**(详见 lib/relay-policy.js)。
|
||
//
|
||
// 对方是 Agent 时它那边的插件也会自动回一封,于是两个模型都以为
|
||
// 「我只要把话说完就行」,实际在持续互相唤醒 —— 生产实测 pi 与 dsh
|
||
// 客套 6 轮直到撞上 hop 上限。这一步在取文本之前判:早退比白做一轮清楚。
|
||
const policy = autoRelayDecision({
|
||
fromHuman: job.data?.from_human === true,
|
||
replyTo: mailContext.replyTo,
|
||
});
|
||
if (!policy.relay) {
|
||
log(`本轮不自动转发:${policy.reason}`);
|
||
return;
|
||
}
|
||
|
||
// 只取 type==='text' 的块(B-5.1 / N-6):thinking 是思考过程,不是结论。
|
||
const text = lastAssistantText(session.messages);
|
||
if (!text) return; // 空文本不发空邮件(B-5.4)
|
||
|
||
// 幂等键用 pi 会话 id + 会话树叶子 id:两者都落盘,重放也是同一个键。
|
||
const relayKey = relayKeyFor(piSessionId, sessionManager.getLeafId?.());
|
||
if (relayedKey === relayKey) return;
|
||
|
||
// 模型这一轮已亲手回过这条线索 → 让位(B-5.3)。
|
||
// 否则收件箱里是两封说同一件事的邮件(生产实测过)。
|
||
if (shouldSkipAutoRelay(explicitSends.get(piSessionId), mailContext.replyTo, mailContext.mailID)) {
|
||
relayedKey = relayKey;
|
||
log(`本轮模型已主动回信 ${mailContext.replyTo},跳过自动转发`);
|
||
return;
|
||
}
|
||
|
||
await client.post('/mail/send', {
|
||
to: mailContext.replyTo,
|
||
subject: replySubject(mailContext.subject),
|
||
body: text,
|
||
reply_to: mailContext.mailID || '',
|
||
// relay + relay_key 走免配额通道(I-2):模型已经把话说完了,桥只是把它
|
||
// 搬到邮件里。对搬运收费会让配额用尽时 Agent 连交代都做不了。
|
||
relay: 'summary',
|
||
relay_key: relayKey,
|
||
});
|
||
relayedKey = relayKey;
|
||
log(`已转发本轮总结给 ${mailContext.replyTo}(${text.length} 字)`);
|
||
}
|
||
|
||
// ─── 主流程 ───
|
||
|
||
async function run() {
|
||
const data = job.data;
|
||
const kind = job.kind;
|
||
|
||
client = new GatewayClient({
|
||
url: job.config.gatewayURL,
|
||
agentName: job.config.agentName,
|
||
agentKey: job.config.agentKey,
|
||
agentSecret: job.config.agentSecret,
|
||
});
|
||
|
||
// ModelRuntime 每个 worker 建一次。allowModelNetwork 保持默认 false:
|
||
// 上报给 Gateway 的是 getAvailable()(有凭证、真能调起来的),
|
||
// 那取决于本机 auth.json,不取决于目录里有多少条。
|
||
modelRuntime = await ModelRuntime.create();
|
||
const runtimeErr = modelRuntime.getError?.();
|
||
if (runtimeErr) log(`模型运行时告警: ${runtimeErr}`);
|
||
|
||
// connect_to_server 在 worker 里换了坐标要让主进程知道:worker 马上就退了,
|
||
// 改在自己身上等于没改。主进程收到后重建 SSE 并写进后续 worker 的 job。
|
||
const mailTools = createMailTools({
|
||
client, log, agentName: job.config.agentName,
|
||
onReconnect: () => send({
|
||
type: 'reconfigure', url: client.baseURL, agentKey: client.agentKey,
|
||
}),
|
||
});
|
||
|
||
const { session, sessionManager, diagnostics, cwd, reused } = await loadSession(mailTools);
|
||
for (const d of diagnostics) log(`扩展诊断: ${d?.message ?? JSON.stringify(d)}`);
|
||
|
||
piSessionId = session.sessionId;
|
||
// 「这条会话是接管来的吗」以 Gateway 给的 platform_session_id 为准,
|
||
// **不能**从 `reused && !job.session.sessionFile` 推断:第二封邮件进同一条
|
||
// 接管会话时主进程给了 sessionFile,那样推断会得出 false,于是命名同步跑起来
|
||
// 把人从补全里选的 slug 冲掉 —— 那正是上一轮修掉的那个 bug。
|
||
const adopted = Boolean(adoptedSessionID(data));
|
||
send({
|
||
type: 'session_opened',
|
||
piSessionId,
|
||
sessionFile: session.sessionFile || '',
|
||
cwd,
|
||
reused,
|
||
});
|
||
|
||
// pi 侧改名(人在 TUI 里 /name)→ 同步给 Gateway。接管会话不同步,理由见
|
||
// relaySummary 的 adopted 参数注释。
|
||
session.subscribe((event) => {
|
||
if (event?.type === 'session_info_changed' && !adopted) {
|
||
syncNaming(session, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`));
|
||
}
|
||
});
|
||
|
||
const prompt = buildMailPrompt({ agentName: job.config.agentName, data, kind, reused });
|
||
|
||
// 轮次超时给得很宽(默认 10 分钟,由主进程的 workerMaxMs 派生)。
|
||
//
|
||
// 原来是 60 秒「超时按成功返回」,因为主进程要腾出手来收下一封邮件;
|
||
// worker 没有这个理由 —— 它只为这封邮件活着,等真结论更准。带工具调用的
|
||
// 一轮跑几分钟很正常,60 秒返回会让转发落在一个还没说完的结论上。
|
||
const turnTimeout = job.config.turnTimeoutMs;
|
||
|
||
if (reused) {
|
||
// 续谈不做模型降级:换模型要换会话,会丢掉整条上下文 —— 而上下文正是
|
||
// 发件人指定这条会话的原因。但失败要可见。
|
||
const outcome = await runTurn(session, prompt, turnTimeout);
|
||
log(`续谈 ${piSessionId}(mail ${data.mail_id}${outcome.queued ? ',已排队' : ''})`);
|
||
if (!outcome.ok) throw new Error(`续谈失败: ${outcome.error}`);
|
||
await relaySummary(session, sessionManager, adopted);
|
||
return;
|
||
}
|
||
|
||
// 按管理员划定的范围逐个尝试(D-3)。
|
||
// 关键点:`prompt()` resolve **不代表模型跑成功了** —— 无凭证的 provider
|
||
// 会让它 reject(实测 `No API key found for amazon-bedrock.`),
|
||
// 而上游报错走 stopReason==='error'。判定交给 classifyTurnOutcome。
|
||
const attempts = modelAttemptOrder(job.config.allowedModels, {
|
||
provider: job.config.replyProvider,
|
||
model: job.config.replyModel,
|
||
});
|
||
const failures = [];
|
||
let live = { session, sessionManager };
|
||
|
||
for (const route of attempts) {
|
||
const label = route ? `${route.provider}/${route.model}` : '(平台默认)';
|
||
if (route) {
|
||
const model = modelRuntime.getModel(route.provider, route.model);
|
||
if (!model) {
|
||
// 目录里根本没有这个路由:同步就能判定,不必起一轮。
|
||
failures.push({ ...route, error: `平台目录里没有 ${label}` });
|
||
log(`模型 ${label} 不存在,跳过`);
|
||
continue;
|
||
}
|
||
// 换模型要换会话:pi 的模型在 createAgentSession 时绑定。
|
||
// 上一次尝试失败的会话没有任何 assistant 消息,丢掉不损失内容。
|
||
try { live.session.dispose?.(); } catch { /* 已经没了 */ }
|
||
const retried = await openSession({
|
||
cwd, modelRuntime, model, customTools: mailTools,
|
||
extension: permissionExtension(),
|
||
});
|
||
live = { session: retried.session, sessionManager: retried.sessionManager };
|
||
piSessionId = retried.session.sessionId;
|
||
send({
|
||
type: 'session_opened',
|
||
piSessionId,
|
||
sessionFile: retried.session.sessionFile || '',
|
||
cwd,
|
||
reused: false,
|
||
});
|
||
}
|
||
const outcome = await runTurn(live.session, prompt, turnTimeout);
|
||
if (outcome.ok) {
|
||
if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`);
|
||
await relaySummary(live.session, live.sessionManager, false);
|
||
return;
|
||
}
|
||
failures.push({ ...(route || {}), error: outcome.error });
|
||
log(`模型 ${label} 失败: ${outcome.error}`);
|
||
}
|
||
|
||
// 全部失败 → 必须回信(B-6):模型一次都没跑起来,会话里没有任何 assistant
|
||
// 消息,自动转发因此什么也不会发 —— 发件人只会看到再无音讯。
|
||
if (kind === 'mail' && data.from_name) {
|
||
try {
|
||
await client.post('/mail/send', {
|
||
to: data.from_name,
|
||
subject: `处理失败: ${data.subject || '(无主题)'}`,
|
||
body: renderFailureReport(failures, data.subject),
|
||
reply_to: data.mail_id || '',
|
||
relay: 'summary',
|
||
relay_key: `model-failure:${data.mail_id || piSessionId}`,
|
||
});
|
||
log(`已回报模型调用失败给 ${data.from_name}`);
|
||
} catch (e) {
|
||
log(`失败回报也发不出去: ${describeError(e)}`);
|
||
}
|
||
}
|
||
// 发完仍要 throw(B-6.4):静默会让这次失败只存在于邮件里,日志上看不出来。
|
||
throw new Error(`范围内 ${failures.length} 个模型全部失败:${failures.map((f) => f.error).join(' | ')}`);
|
||
}
|
||
|
||
// ─── 入口 ───
|
||
|
||
process.on('message', (msg) => {
|
||
if (msg?.type === 'job') {
|
||
if (job) return; // 一个 worker 只接一封
|
||
job = msg;
|
||
mailContext = {
|
||
replyTo: msg.data?.from_name || '',
|
||
subject: msg.data?.subject || '',
|
||
mailID: msg.data?.mail_id || '',
|
||
permissionMode: msg.data?.permission_mode || 'workspace',
|
||
};
|
||
lastSyncedName = msg.lastSyncedName || '';
|
||
for (const t of msg.grants || []) grants.add(t);
|
||
run()
|
||
.then(() => finish({ ok: true, error: '' }))
|
||
.catch((e) => finish({ ok: false, error: describeError(e) }));
|
||
return;
|
||
}
|
||
if (msg?.type === 'permission_decision') {
|
||
const resolve = pending.get(msg.relayKey);
|
||
if (!resolve) return;
|
||
pending.delete(msg.relayKey);
|
||
resolve(String(msg.decision || '拒绝'));
|
||
return;
|
||
}
|
||
if (msg?.type === 'shutdown') {
|
||
// fail closed(B-9.2 / N-9):唤醒所有未决询问,让 pi 侧那些 await 返回。
|
||
// 不唤醒的话会话挂死;而默认放行一个没人批准的危险操作更糟。
|
||
for (const [key, resolve] of pending) {
|
||
log(`未决权限 ${key} fail closed`);
|
||
resolve('shutdown');
|
||
}
|
||
pending.clear();
|
||
}
|
||
});
|
||
|
||
function finish(result) {
|
||
if (finished) return;
|
||
finished = true;
|
||
send({ type: 'done', ...result });
|
||
// **一律 exit 0**:`done` 已经把成败说清楚了,用退出码再说一遍会让主进程
|
||
// 把「模型全部失败」这种已处理的结果也打成「worker 异常退出」。
|
||
// 非零退出码留给真正的崩溃(没来得及发 done 的那种)。
|
||
//
|
||
// 给 IPC 一个 tick 把 done 送出去再退:未 flush 的消息会丢,而主进程靠它判成败。
|
||
// unref 不影响触发(IPC 通道本身在保持事件循环存活),只是不由这个计时器兜着。
|
||
const t = setTimeout(() => process.exit(0), 50);
|
||
if (typeof t.unref === 'function') t.unref();
|
||
}
|
||
|
||
// 未捕获异常也要回报:静默退出会让主进程只看到 exit code,
|
||
// 而那条邮件的失败原因就此丢失。
|
||
process.on('uncaughtException', (e) => finish({ ok: false, error: `未捕获异常: ${describeError(e)}` }));
|
||
process.on('unhandledRejection', (e) => finish({ ok: false, error: `未处理拒绝: ${describeError(e)}` }));
|
||
|
||
send({ type: 'ready' });
|