/** * pi 会话池 —— 每条 AgentMail 会话对应一条 pi 会话。 * * 为什么桥必须自己持有 pi 会话(而不是写成一个 pi 扩展): * 扩展被加载进**一条已经存在的**会话里,cwd 由启动 pi 的人决定;而 B-3.1 要求 * 每封邮件的 to_workspace 成为会话 cwd。扩展做不到「按邮件新开一条 cwd 不同的 * 会话」,所以桥是一个常驻进程(C-7),用 SDK 的 createAgentSession 起会话。 * * 每条会话一套 SettingsManager / ResourceLoader / SessionManager:它们都按 cwd * 解析项目级配置(.pi/、skills、prompts),共用一份会把 A 项目的配置带进 B 项目。 */ import { createAgentSession, SessionManager, SettingsManager, DefaultResourceLoader, getAgentDir } from '@earendil-works/pi-coding-agent'; /** * 起一条 pi 会话。 * * @param {object} opts * @param {string} opts.cwd 会话工作目录(已由 resolveWorkspaceCwd 校验过存在) * @param {any} opts.modelRuntime 共享的 ModelRuntime(建一次很贵,池外传进来) * @param {any} [opts.model] 指定模型;省略则用 settings 里的默认 * @param {any[]} opts.customTools 邮件工具(send_mail / read_inbox / …) * @param {(pi: any) => void} [opts.extension] 内联扩展工厂,用来挂 tool_call 权限钩子 * @returns {Promise<{session: any, sessionManager: any, diagnostics: any[]}>} */ export async function openSession({ cwd, modelRuntime, model, customTools, extension, sessionFile }) { const agentDir = getAgentDir(); const settingsManager = SettingsManager.create(cwd, agentDir); const resourceLoader = new DefaultResourceLoader({ cwd, agentDir, settingsManager, // 关掉磁盘上的全局扩展。两个理由: // 1. 本机的 pi-a2a / pi-acp 在加载时 listen 固定端口(12010/12011), // 守护进程里加载会 EADDRINUSE,把整条会话拖死。 // 2. 桥起的会话是给邮件用的,不该继承人类交互用的那套扩展(TUI 命令、 // 快捷键、状态栏都没有意义)。 // 邮件工具走 customTools,权限钩子走下面的 extensionFactories。 noExtensions: true, extensionFactories: extension ? [{ name: 'agentmail-bridge', factory: extension }] : [], }); await resourceLoader.reload(); // sessionFile 非空 = **接管一条磁盘上已经存在的会话**(人在 TUI 里开的那种)。 // // `SessionManager.open` 把整条会话装回内存(历史消息、分支、标签都在), // 之后 prompt 就是在那条对话后面接着谈 —— 人在 TUI 里再打开它能看到 // 邮件带来的这一轮。TUI 与邮箱是同一个 Agent 的两个入口。 // // # 双写风险与它的边界 // // pi 没有任何锁机制(SDK 里 flock/lockfile 命中为 0),它假定「一个文件 // 一个持有者」。写入本身是纯 append(`_persist` → `appendFileSync`),所以 // 两个持有者不会把文件截断;坏的是**各自的内存索引**:对方追加的行自己看不见, // 于是算出的 parentId 指向一个对方不知道的 entry,会话树分叉。 // // 取舍是「短暂持有」:open → 跑一轮 → 丢弃这个 manager(调用方不缓存它)。 // 窗口是一轮对话的时长。人正好在那一刻也在 TUI 里发消息仍会分叉 —— // 但那需要两边同时动手,而分叉的后果是历史看起来少了一段,不是数据损坏。 // // cwd 用会话 header 里的(open 的第三参不传即取 header),不是外面传进来的: // 会话的工作目录在它创建时就定了,传一个不同的只会让项目级配置错位。 const sessionManager = sessionFile ? SessionManager.open(sessionFile) : SessionManager.create(cwd); const created = await createAgentSession({ cwd, agentDir, modelRuntime, // model 为 undefined 时 SDK 用 settings 里的默认模型,正好对应 // modelAttemptOrder 里那个 `undefined`(= 不指定、交给平台)。 ...(model ? { model } : {}), sessionManager, settingsManager, resourceLoader, customTools, }); return { session: created.session, sessionManager, diagnostics: created.extensionsResult?.diagnostics ?? [], }; } /** * 跑一轮并等到真正的结论(C-4 / D-3)。 * * `session.prompt()` 的 promise 在**这一轮彻底结束**时才 resolve,所以不需要 * 额外订阅 agent_end 去等。但它 resolve 了**不代表模型跑成功了** —— * 判定交给 classifyTurnOutcome(三条互不重叠的失败信号,见那里的注释)。 * * 60 秒超时算成功(与另两个插件同一取舍):长任务很正常,把它判成失败会 * 换模型重跑一遍,等于同一封邮件跑两次。超时只是「不再等着上报结论」, * 会话仍在跑,轮次结束后 agent_end 会照常触发自动转发。 * * 会话正在跑时走排队(返回 queued),**不能**在那种情况下判结论: * prompt 排完队就 resolve,此时 session.messages 里最后一条是**上一轮**的, * 拿它判定会把上一轮的成败当成这一轮的。 * * @param {any} session * @param {string} promptText * @param {number} timeoutMs * @returns {Promise<{ok: boolean, error: string, aborted: boolean, timedOut: boolean, queued: boolean}>} */ export async function runTurn(session, promptText, timeoutMs = 60_000) { const { classifyTurnOutcome } = await import('./turn.mjs'); // 排队分支:模型还在说话时又来一封邮件。 // // streamingBehavior 必选,缺了 prompt 直接抛 // "Agent is already processing. Specify streamingBehavior…"。 // 取 followUp 而不是 steer:steer 会把当前这一轮打断, // 而当前这一轮正在处理**上一封邮件** —— 那封邮件的发件人也在等回信。 if (session.isStreaming) { await session.prompt(promptText, { streamingBehavior: 'followUp' }); return { ok: true, error: '', aborted: false, timedOut: false, queued: true }; } let timer = null; const timeout = new Promise((resolve) => { timer = setTimeout( () => resolve({ ok: true, error: '', aborted: false, timedOut: true, queued: false }), timeoutMs, ); }); const run = session.prompt(promptText) .then(() => ({ ...classifyTurnOutcome({ messages: session.messages }), timedOut: false, queued: false })) .catch((e) => ({ ...classifyTurnOutcome({ error: e }), timedOut: false, queued: false })); try { return await Promise.race([run, timeout]); } finally { if (timer) clearTimeout(timer); } } /** * 续谈:往一条已经存在的会话里追加一轮。 * * 这就是 `runTurn` —— 不需要第二个函数。 * * **不能**用 `session.followUp()`:那个方法只往 followUpQueue 里塞消息, * 队列**只在运行中的轮次末尾**被 drain(pi-agent-core/agent.js 的 run 循环, * 以及 `continue()`)。会话空闲时(上一轮早已结束)塞进去的消息永远没人取, * 于是这封邮件既没有回信也没有报错 —— 实测踩过:日志打了「续谈」, * 收件箱里只有来信没有回复。 * * `runTurn` 按 `isStreaming` 分流,两种状态都正确: * - 空闲 → `prompt()` 直接起一轮 * - 正在跑 → `prompt(text, {streamingBehavior:'followUp'})` 排到当轮之后 */ export { runTurn as followUpTurn };