Files
MailUI4Agents/plugins/pi-mail-bridge/src/session-pool.mjs
JianFeeeee 255c799a40 feat(adopt): 邮件可投进平台上已存在的会话(TUI 与邮箱同一入口)
人在平台界面(pi TUI / opencode / DSH GUI)里开的会话,此前无法被邮件投进去。
补全早就把它们列为候选(agent_platform_sessions 镜像,插件心跳上报),
但投递侧的 FindNamedSessionFor 只查 sessions 表 —— 选中后只能得到 404。
候选列表在承诺一件做不到的事。

TUI 与邮箱是同一个 Agent 的两个入口,不是两套隔离的世界。

## Gateway

sessions 表加 platform_id 列 + 部分索引。resolveTarget 的 SessionNamed 分支
本侧查不到时再查镜像,命中则「接管」:本侧建一条会话并绑定 platform_id,
之后每次投递都在 SSE 事件里带 platform_session_id。

- FindPlatformSession(agent, slug, workspace) 查镜像
- FindSessionByPlatformID 防重复接管(一条平台会话只能被接管一次,
  否则同一条对话在邮箱里裂成多条互不相干的线索)
- AdoptPlatformSession 建会话 + 绑定 + 别名复用平台 slug(撞名自动加后缀)
- PlatformIDOf 供 notifyRecipients 读

三处语义决定:
- workspace 以平台会话为准(它的 cwd 创建时就定了)。地址 path 位不同则不命中,
  否则邮件会投进另一个项目的会话
- 主题优先用平台侧标题(它代表整条对话在谈什么,也是补全里显示的)
- 接管计入 AllowNewSession 速率限制 —— 镜像里可能有几百条 slug,
  不计的话它是绕过限流的后门

## 插件

字段解析与失败话术抽成共用模块 lib/adopt.js(三方逐字节相同 + 进同源校验):
字段名各写一遍时少个下划线就静默退化成「每封邮件新开一条」,而那个错误不抛异常。

- opencode:session.get 确认存在 → 照常 promptAsync(服务端持有会话,单一写者)
- DSH:复用 startAgent 的 resume 分支,会话 id 换成平台自己那个;
  界面上正开着时直接 followup(两个 handle 会各自写日志,replay 过不去)
- pi:SessionManager.open(file) → 跑一轮 → dispose,不放进长期缓存

pi 必须短暂持有:SDK 无任何锁机制(flock/lockfile 命中 0),活着的
SessionManager 不 watch 文件 —— 外部追加的行看不见,算出的 parentId 指向
对方不知道的 entry,会话树分叉。写入是纯 append 所以文件不会坏。
配套三处:isStreaming 时不释放(否则杀掉排队中的下一封)、兜底计时器
(轮次超时 ×2,unref)、接管会话跳过命名同步。

最后一条是实测撞出来的:别名撞名时 Gateway 加后缀,而定稿别名又回写进 pi
会话文件 → 下次心跳上报的 slug 变成带后缀那个,人从补全里选的名字凭空消失。
opencode/DSH 无此环(它们的 slug 只读不写)。

接管后必须加入 mailDriven 集合,否则邮件投进去了却永远没有回音。

## 迁移顺序

idx_sessions_platform 不能写在 init_sqlite.sql 里:那个脚本在
addMissingColumns 之前执行,而已部署的库里 sessions 表已存在
(CREATE TABLE IF NOT EXISTS 不补列)→ 索引建在不存在的列上,
整个迁移中断、服务起不来(生产实测)。依赖补出来的列的索引一律放
migrate.go 的 sqliteAddIndexes。PG 侧用 ALTER TABLE ADD COLUMN IF NOT EXISTS。

## 生产验证

- pi × 2(agent-only-chain / mail-probe-alias)、opencode(glowing-moon)、
  dsh(查看工程与插件适配指南)四条链路接管成功
- dsh 那次回信准确说出了界面上聊过的内容 → 上下文确实装回来了
- 第二封复用同一条本侧会话,平台侧无新增改名条目
- 回归:opencode 普通 .new + 别名续谈 + used_rounds=0(免配额通道未受影响)

## 其他

pi-mail-bridge 补 systemd 单元(此前是 setsid 裸进程,重启机器不会拉起):
陈锁清理 ExecStartPre、MemoryMax=4G、TimeoutStopSec=10。
配置目录必须与 opencode 分开(共用会让后起的读到对方密钥或撞单实例锁)。

PLUGIN-CONTRACT.md 加 B-3.7 / B-3.8 + new_mail 字段表 + 检查清单验收项。

测试:repo +10 例(adopt_test.go);三插件各 +7 例(adopt.test.mjs)
2026-09-04 11:14:44 +08:00

159 lines
7.3 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.

/**
* 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 而不是 steersteer 会把当前这一轮打断,
// 而当前这一轮正在处理**上一封邮件** —— 那封邮件的发件人也在等回信。
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 里塞消息,
* 队列**只在运行中的轮次末尾**被 drainpi-agent-core/agent.js 的 run 循环,
* 以及 `continue()`)。会话空闲时(上一轮早已结束)塞进去的消息永远没人取,
* 于是这封邮件既没有回信也没有报错 —— 实测踩过:日志打了「续谈」,
* 收件箱里只有来信没有回复。
*
* `runTurn` 按 `isStreaming` 分流,两种状态都正确:
* - 空闲 → `prompt()` 直接起一轮
* - 正在跑 → `prompt(text, {streamingBehavior:'followUp'})` 排到当轮之后
*/
export { runTurn as followUpTurn };