#!/usr/bin/env node /** * AgentMail ↔ pi 桥(pi-mail-bridge) * * 形态是**常驻守护进程**,不是 pi 扩展。原因见 src/session-pool.mjs 顶部: * 扩展被加载进一条已存在的会话,cwd 由启动 pi 的人决定;而 B-3.1 要求每封邮件的 * to_workspace 成为会话 cwd。 * * # 进程结构 * * 主进程**只做 I/O 与调度**:SSE、心跳、去重、把邮件派给子进程。模型工作全部 * 下到 `src/worker.mjs`(一封邮件一个进程,跑完就退),由 `src/pool.mjs` 调度。 * * 这不是为了并行度,是为了**不阻塞事件循环**。pi 的会话装载是同步的: * `SessionManager.open()` 走 `openSync` + `readSync` 循环把整个 `.jsonl` 读进内存 * 并逐行 JSON.parse。实测本机最大那条会话 23MB,`open` 一次阻塞事件循环 118ms; * 模型跑起来之后 SDK 内部还有更多同步工作。原来这些都在主线程上 —— SSE 读循环 * 在那期间完全停住,后续邮件卡在 TCP 缓冲区,久到 Gateway 认为连接死了, * 重连又触发重放。实测同样的活在 fork 出的子进程里跑,主进程阻塞 0ms。 * * 权限询问期间的挂起也随之只影响那一个 worker:原来 `await new Promise(...)` * 等人做决定,整座桥在那段时间不再收信。 * * 契约实现对照(docs/PLUGIN-CONTRACT.md): * B-1 启动 → main() * B-2 心跳 → beat(),30 秒 * B-3 new_mail → pool.submit()(投递本体在 worker.mjs) * B-4 决策 → handlePermissionDecision() * B-5 转发 → worker.mjs 的 relaySummary() * B-6 失败回信 → worker.mjs 末尾的 renderFailureReport * B-7 补拉 → catchUp() * B-8 权限 → worker.mjs 的 permissionExtension() * B-9 关停 → shutdown() */ import { mkdirSync, openSync, closeSync, unlinkSync, readFileSync, writeFileSync } from 'node:fs'; import { homedir } from 'node:os'; import { join } from 'node:path'; import { ModelRuntime, getAgentDir } from '@earendil-works/pi-coding-agent'; import { GatewayClient, readLocalKey, generateLocalKey, saveConfig } from './gateway.mjs'; import { createWorkerPool } from './pool.mjs'; import { createSessionScanner } from './session-scan.mjs'; import { describeError } from './turn.mjs'; import { BoundedSet, MAX_TRACKED_MAILS } from '../lib/bounded.js'; import { snapshotPiModels } from '../lib/model-scope.js'; import { snapshotPiSessions } from '../lib/session-snapshot.js'; import { selectCatchup } from '../lib/catchup.js'; // ─── 配置 ─── const GATEWAY_URL = process.env.AGENTMAIL_GATEWAY_URL || 'http://127.0.0.1:8180'; const AGENT_NAME = process.env.AGENTMAIL_AGENT_NAME || 'pi'; const AGENT_SECRET = process.env.AGENTMAIL_AGENT_SECRET || ''; const REPLY_PROVIDER = process.env.AGENTMAIL_REPLY_PROVIDER || ''; const REPLY_MODEL = process.env.AGENTMAIL_REPLY_MODEL || ''; /** * 单轮超时。 * * 从 60 秒放宽到 10 分钟:60 秒那个数字是「主进程要腾出手来收下一封邮件」的 * 产物 —— 超时按成功返回,好让 deliverMail 早点结束。worker 没有这个理由, * 它只为这封邮件活着,等真结论更准。带工具调用的一轮跑几分钟很正常, * 60 秒返回会让转发落在一个还没说完的结论上。 */ const TURN_TIMEOUT_MS = Number(process.env.AGENTMAIL_TURN_TIMEOUT_MS || 600_000); /** * 并发上限。 * * 每个 worker 约 140MB RSS(实测),并且每个都会对上游 provider 发请求。 * 3 是内存与吞吐的折中;同一条会话无论如何都是串行的(见 pool.mjs)。 */ const MAX_WORKERS = Number(process.env.AGENTMAIL_MAX_WORKERS || 3); /** * worker 硬超时。 * * 比轮次超时留出余量:正常情况下 worker 自己会在轮次超时后收尾退出, * 这个数字兜的是「连收尾都没做」(进程卡死、权限等不到决策而决策事件也丢了)。 * 到点 SIGKILL —— 否则那条会话的后续邮件永远排队。 */ const WORKER_MAX_MS = Number(process.env.AGENTMAIL_WORKER_MAX_MS || TURN_TIMEOUT_MS + 120_000); const LOCK_FILE = join(process.env.AGENTMAIL_CONFIG_DIR || join(homedir(), '.agentmail'), 'pi-bridge.lock'); /** 日志一律 console.error:它一定进 journalctl(契约 9.8)。 */ const log = (...args) => console.error('[pi-mail-bridge]', ...args); // ─── 进程内状态 ─── // // 主进程只留「调度需要的」那几样,全部在内存(重启即丢,契约第六节的已知取舍)。 // 会话映射、权限授权、命名指纹都下沉到 pool 里按邮件会话存 —— 主进程不再持有 // AgentSession 对象(那东西跨不了进程边界)。 // 已投过的 mail_id(SSE 与补拉共用,B-7.3)。 // // 有界:桥是守护进程,跑几十天下来这里会攒下每一封处理过的邮件 id 而永远 // 没有出口。淘汰是安全的 —— 它防的两种重复(心跳与 SSE 建连之间的窗口、 // SSE 断线重放)都发生在秒到分钟级,几千封之前的 id 不可能再来。 const deliveredMails = new BoundedSet(MAX_TRACKED_MAILS); let allowedModels = []; let modelRuntime = null; let client = null; let pool = null; let sessionScanner = null; let heartbeatTimer = null; let shuttingDown = false; // ─── 单实例锁 ─── // // 两个桥同时跑的后果不是「慢一点」而是错的:两条 SSE 各收到同一封邮件, // 各起一条 pi 会话,发件人收到两封回信;而 deliveredMails 在各自内存里,去重不了。 function acquireLock() { mkdirSync(join(LOCK_FILE, '..'), { recursive: true, mode: 0o700 }); try { // O_EXCL 原子创建。存在则说明有别的实例(或上次崩溃留下的陈锁)。 const fd = openSync(LOCK_FILE, 'wx'); writeFileSync(fd, String(process.pid)); closeSync(fd); return true; } catch (e) { if (e?.code !== 'EEXIST') throw e; } // 陈锁判定:文件里的 pid 还活着吗 let pid = 0; try { pid = Number(readFileSync(LOCK_FILE, 'utf8').trim()); } catch { /* 读不到当陈锁 */ } if (pid > 0) { try { // signal 0 只探测存在性,不真的发信号 process.kill(pid, 0); log(`已有实例在运行(pid ${pid}),本进程退出。`); return false; } catch { // ESRCH:进程没了,是陈锁 } } log(`清理陈锁 ${LOCK_FILE}(原 pid ${pid || '未知'} 已不存在)`); try { unlinkSync(LOCK_FILE); } catch { /* 竞态下别人清掉了也行 */ } return acquireLock(); } function releaseLock() { try { // 只删自己的锁:pid 不符说明这把锁已被别的实例接管 if (Number(readFileSync(LOCK_FILE, 'utf8').trim()) === process.pid) unlinkSync(LOCK_FILE); } catch { /* 已经没了 */ } } // ─── 权限决策回来(B-4)─── /** * 把决策路由给发起询问的那个 worker。 * * 找不到 worker 有两种情形,都不该新开会话(B-4.3): * - 桥重启了:那次工具调用早已随进程消失。但人刚刚点了「同意」—— * 什么都不做的话人以为自己批准了、Agent 却毫无反应,所以退化为把决策 * 当一封通知投进原会话(B-4.2)。 * - 那条邮件会话从没被处理过:连通知都无处可投,只能记一行日志。 */ function handlePermissionDecision(data) { const relayKey = data.relay_key || ''; if (relayKey && pool.routePermission(relayKey, String(data.decision || '拒绝'))) { log(`权限 ${relayKey} 决策 ${data.decision}(决策人 ${data.decided_by || '?'})已转交 worker`); return; } if (!data.session_id || !pool.hasSession(data.session_id)) { // **不得凭空新开会话**(B-4.3) log(`权限决策 ${relayKey} 无对应会话,忽略`); return; } log(`权限 ${relayKey} 无挂起项,退化为通知投递`); pool.submit('permission', data); } // ─── 心跳(B-2)─── async function reportSessions() { try { // **不用 `SessionManager.listAll()`**:它为了拿 id/cwd/name/modified 四个 // 字段,把 ~/.pi/agent/sessions 下每个 .jsonl 的每一行都读进来并 JSON.parse, // 还把所有消息正文拼成一个 allMessagesText 大字符串。本机实测(115 个文件 / // 145MB)单次 1431ms、堆里瞬时 240MB —— 而这 282MB 每 30 秒分配一次随即 // 变成垃圾,且那 1.4 秒是同步解析,跑在事件循环上(SSE 读循环那期间停着)。 // // sessionScanner 只读 header 的首行 + 增量扫尾部找 session_info: // 稳态下未变化的文件一个字节都不读(实测 3ms / 0 字节)。 const all = await sessionScanner.scan(); const driven = pool.mailDrivenIDs(); return snapshotPiSessions(all, (id) => driven.has(id)); } catch (e) { // 拉不到就**省略字段**而不是传 [](N-7 / W-3): // 空数组的语义是「平台确实一条会话都没有」,会把服务端镜像抹掉。 log(`会话列表读取失败: ${describeError(e)}`); return undefined; } } async function reportModels() { try { // getAvailable 而不是 getModels:后者本机有 1221 条,其中真能调起来的只有 1 条。 // 上报目录的全部意义就是让管理员别选中一个注定失败的路由。 const available = await modelRuntime.getAvailable(); return snapshotPiModels(available); } catch (e) { log(`模型目录读取失败: ${describeError(e)}`); return undefined; } } /** * 补投离线期间积压的未读邮件(B-7)。 * * SSE 只推连上之后的事件,插件重启前发来的邮件不会再推一次。 * * 与旧版的差别:**不再 await 每一封**。旧版串行是因为「每封都要起一轮模型, * 并发放出去等于对上游打 N 个并发请求」—— 那个约束现在由 pool 的 maxWorkers * 承担,而且它比串行更好:同一条会话仍然串行,不同会话可以并行。 */ async function catchUp(pending) { if (!pending) return; try { const box = await client.get('/mail/inbox?status=unread&limit=20'); const tasks = selectCatchup(box?.mails ?? box, deliveredMails); if (!tasks.length) return; log(`补投 ${tasks.length} 封离线期间的邮件(共 ${pending} 封未读)`); for (const ev of tasks) { if (deliveredMails.has(ev.mail_id)) continue; // 逐封再查(B-7.6) deliveredMails.add(ev.mail_id); pool.submit('mail', ev); } } catch (e) { log(`补投失败: ${describeError(e)}`); } } // ─── 启动 / 关停 ─── async function main() { if (!acquireLock()) process.exit(0); // B-1.1:环境变量 → ~/.agentmail/agent.key → 本地生成并打印全文 let agentKey = process.env.AGENTMAIL_AGENT_KEY || readLocalKey(); if (!agentKey && !AGENT_SECRET) agentKey = generateLocalKey(log); client = new GatewayClient({ url: GATEWAY_URL, agentName: AGENT_NAME, agentKey, agentSecret: AGENT_SECRET, }); // ModelRuntime 主进程也要一个:心跳的 reportModels 用它。worker 各自再建 // 一个(跨进程传不了),代价是每个 worker 多 ~20ms(实测 11–23ms)。 // // allowModelNetwork 保持默认的 false:桥启动时不去网上拉模型目录。 // 拉了也没用 —— 上报给 Gateway 的是 getAvailable()(有凭证、真能调起来的), // 而那取决于本机 auth.json,不取决于目录里有多少条。开着只会让 // 启动多等一个网络往返,而且断网时启动路径上多一个可失败点。 modelRuntime = await ModelRuntime.create(); const runtimeErr = modelRuntime.getError?.(); if (runtimeErr) log(`模型运行时告警: ${runtimeErr}`); // 会话目录扫描器。**必须建一次并复用** —— 它的省内存全靠跨拍存活的 // size 缓存(稳态下未变化的文件一个字节都不读)。每拍新建一个等于 // 每拍都冷启动,退回 listAll 那种全量读的开销。 // // 路径自己拼而不是 import getSessionsDir:SDK 只导出 getAgentDir, // getSessionsDir 是内部函数(dist/config.js 里 `join(getAgentDir(), "sessions")`)。 sessionScanner = createSessionScanner({ sessionsDir: join(getAgentDir(), 'sessions'), }); // 工作进程池。config() 每次派活时取一次 —— allowedModels 随心跳变, // 取快照会让 worker 用上一轮的模型范围。 pool = createWorkerPool({ log, config: () => ({ gatewayURL: client.baseURL, agentName: AGENT_NAME, agentKey: client.agentKey, agentSecret: AGENT_SECRET, allowedModels, replyProvider: REPLY_PROVIDER, replyModel: REPLY_MODEL, turnTimeoutMs: TURN_TIMEOUT_MS, }), // worker 里 connect_to_server 换了坐标:worker 马上就退了,改在它自己身上 // 等于没改。主进程据此重建 SSE,后续 worker 的 job 也会带上新坐标。 onReconfigure: (url, key) => { if (!client.reconfigure({ url, agentKey: key })) return; log(`Gateway 坐标已更新为 ${client.baseURL},重建 SSE`); client.stopSSE(); client.startSSE(handleSSEEvent, log); }, maxWorkers: MAX_WORKERS, workerMaxMs: WORKER_MAX_MS, }); // 主进程不跑模型,因此**不建**邮件工具:工具是给模型调的,而这里没有会话。 // (工具 schema 的约束由 test/tool-schema.test.mjs 直接验证 createMailTools, // 不需要在这里建一份没人用的副本。) // // connect_to_server 换坐标的闭环在 pool 的 onReconfigure 里 —— 那个工具 // 跑在 worker 里,worker 把新坐标回报给主进程,主进程据此重建 SSE。 try { await client.register(); // B-1.2 saveConfig({ gateway_url: GATEWAY_URL, agent_name: AGENT_NAME, registered_at: new Date().toISOString() }); log(`已接入 ${GATEWAY_URL},身份 ${AGENT_NAME}(${agentKey ? '密钥认证' : 'name/secret 认证'})。`); } catch (e) { // 密钥未登记时这里报「密钥无效」—— 必须说清该做什么, // 否则用户只看到一句 401,不知道要拿密钥去后台登记。 log(`注册失败: ${describeError(e)}`); if (agentKey) log(`若提示密钥无效,请让管理员在 AgentMail 后台登记这把密钥。`); } let caughtUp = false; const beat = async () => { const [platform_sessions, models] = await Promise.all([reportSessions(), reportModels()]); const body = {}; if (platform_sessions) body.platform_sessions = platform_sessions; if (models) body.models = models; try { const res = await client.post('/agent/heartbeat', body); if (Array.isArray(res?.allowed_models)) allowedModels = res.allowed_models; // B-2.2 if (!caughtUp) { // B-7.1:只在首个成功心跳后补一次 caughtUp = true; await catchUp(res?.pending_mails); } } catch { // B-2.1:心跳失败不重试不报错。真连不上时 Gateway 会把它判成离线, // 那才是可见的信号;桥自己打一串错误日志只会淹掉真正的问题。 } }; await beat(); // B-1.3:不等第一个 30 秒周期 heartbeatTimer = setInterval(beat, 30_000); // B-1.5 client.startSSE(handleSSEEvent, log); for (const sig of ['SIGINT', 'SIGTERM']) process.on(sig, () => shutdown(sig)); } /** * SSE 事件分派。 * * 这个函数**必须保持廉价**:它跑在读循环上。派活给 pool 是同步的(fork 是 * 异步的,pool.submit 只是入队),所以读循环不会因为一封邮件停下。 * * 提成命名函数是因为 connect_to_server 换地址后要用同一个处理器重建长连 —— * 内联箭头函数在那里拿不到,只能复制一遍,而复制出来的两份迟早会分叉。 */ function handleSSEEvent(type, data) { if (type === 'permission_decision') { handlePermissionDecision(data); return; } if (type === 'session_archived') { // 会话归档 = 那条会话再也不会收信(别名 404),pool 里的 sessionState // 可以确定性地清掉,不必等上限淘汰去猜。 if (pool?.forget(data?.session_id || '')) { log(`会话 ${data.session_id} 已归档,清除本地状态`); } return; } if (type !== 'new_mail') return; if (data?.role && data.role !== 'to' && data.role !== 'cc') return; const id = data?.mail_id; if (!id || deliveredMails.has(id)) return; // B-3 第 1 步:去重 deliveredMails.add(id); pool.submit('mail', data); } function shutdown(reason) { if (shuttingDown) return; shuttingDown = true; log(`收到 ${reason},关停中…`); if (heartbeatTimer) clearInterval(heartbeatTimer); // B-9.1 client?.stopSSE(); // pool.stop 先给每个 worker 发 shutdown(让它把未决权限询问 fail closed, // B-9.2 / N-9),再给 2 秒自己退,然后 SIGKILL。 // // 不直接杀:pi 侧那些 await 不会返回,而 worker 里可能正握着会话文件。 pool?.stop(); releaseLock(); // 留出 pool.stop 的宽限窗口再退:主进程先死会让子进程变成孤儿 // (systemd 的 KillMode 会兜住,但那时 fail closed 已经来不及做了)。 setTimeout(() => process.exit(0), 2500).unref?.(); // B-9.3:不发「插件下线」通知邮件 } main().catch((e) => { log(`启动失败: ${describeError(e)}`); releaseLock(); process.exit(1); });