feat: agent 邮件寻址能力全面补齐 + .new 别名替换
## 别名替换(让 .new 邮件可寻址)
repo/autoalias.go: AutoAliasFor + EnsureSessionAlias
- .new 建完会话立刻给别名(形如 dsh-重构导入路径)
- 名字与主题都要:只用主题跨 Agent 撞名,只用名字看不出聊什么
- sanitizeAliasPart 只留 unicode.IsLetter/IsDigit,其余折 -
- 撞名追加 -2/-3,全占用退 session-<uuid前8位>
- 不复用 SyncSessionAlias:那个假定已存在且跳过 manual
- 条件写入 WHERE alias IS NULL OR '',并发安全
- resolveTarget 的 .new 与默认会话两条路径都调
notifyRecipients 加三个字段(每个收件方拿到自己那个地址的版本):
- session_alias / reply_address / self_address
- 别名为空时退回省略 session 位,绝不写 new
FormatAddress(name,path,session) 空 path 也必须留 @ 与 .
## Agent 侧寻址发现(五个只读端点)
handler/agent_discovery.go:
- /agent/contacts + /agent/contacts/suggest(三段式补全)
- /agent/mail/{id} + /agent/mail/{id}/thread
- /agent/sessions/{id}/participants
- 不复用人类路由:scope 不同、审计需求不同
- 一律只读:归档/改名/权限决策仍只有人能做
repo/participants.go: SessionParticipants 逐封扫 from/to/cc
- Roles 用集合、MailCount 只数发信(0=还没开口的人)
- 发件人 path 不取 from_workspace(那列存的是 Agent 名)
repo.SuggestPaths 重写:mails.to_workspace(按 MAX(created_at) 倒序)
+ agents.workspaces 并集。原只读 workspaces,官方插件传 [] 永远空
## 共用模块(三插件逐字节相同)
lib/addressing.js: formatAddress/roleOf/replyAddressFor/selfAddressFor/participantsOfMail
lib/discovery.js: renderNameSuggestions/renderPathSuggestions/renderSessionSuggestions/
renderParticipants/renderContacts/renderThread
lib/inbox-format.js: renderMail 新增收件人/身份/可投递地址三段
- selfName 参数(兼容旧调用不传的情况)
check-shared-libs.sh 纳入 addressing + discovery
## 插件侧
opencode: suggest_address + list_contacts + session_participants + read_thread + read_mail
dsh: 同上 + forward_mail(此前只有 opencode 有)+ upload_attachment 改真 multipart
pi: 同上(createMailTools 加 agentName 参数)
dsh: ctx.agents.create id collision 改为 readSession 探测后 resume
dsh: 关键路径日志改 console.error(ctx.logger 不进 journalctl)
## 测试
repo: autoalias_test.go 11 + participants_test.go 7 = 18 例
plugins: addressing.test 17 + discovery.test 23 + inbox-format.test 31 = 71 例
go test ./... + npm test(opencode 155 + dsh 173 + pi 199)全绿
端到端验证:admin 发 dsh@....new 抄送 opencode@....new
→ dsh 用 session_participants 取到地址 → send_mail 给 opencode
→ 地址取自工具返回值(.crisp-planet),未手工拼写
This commit is contained in:
693
plugins/pi-mail-bridge/src/index.mjs
Normal file
693
plugins/pi-mail-bridge/src/index.mjs
Normal file
@ -0,0 +1,693 @@
|
||||
#!/usr/bin/env node
|
||||
/**
|
||||
* AgentMail ↔ pi 桥(pi-mail-bridge)
|
||||
*
|
||||
* 形态是**常驻守护进程**,不是 pi 扩展。原因见 src/session-pool.mjs 顶部:
|
||||
* 扩展被加载进一条已存在的会话,cwd 由启动 pi 的人决定;而 B-3.1 要求每封邮件的
|
||||
* to_workspace 成为会话 cwd。桥用 SDK 的 createAgentSession 按邮件起会话,
|
||||
* 一个进程里并存多条不同 cwd 的会话(实测可行)。
|
||||
*
|
||||
* 契约实现对照(docs/PLUGIN-CONTRACT.md):
|
||||
* B-1 启动 → main()
|
||||
* B-2 心跳 → beat(),30 秒
|
||||
* B-3 new_mail → deliverMail()
|
||||
* B-4 决策 → handlePermissionDecision()
|
||||
* B-5 转发 → relaySummary(),挂在 agent_end 上
|
||||
* B-6 失败回信 → deliverMail() 末尾的 renderFailureReport
|
||||
* B-7 补拉 → catchUp()
|
||||
* B-8 权限 → permissionExtension() 的 tool_call 钩子
|
||||
* 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 } from '@earendil-works/pi-coding-agent';
|
||||
|
||||
import { GatewayClient, readLocalKey, generateLocalKey, saveConfig } 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, snapshotPiModels } from '../lib/model-scope.js';
|
||||
import { snapshotPiSessions } from '../lib/session-snapshot.js';
|
||||
import { selectCatchup } from '../lib/catchup.js';
|
||||
import { explicitSends, shouldSkipAutoRelay } from '../lib/relay-dedup.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 || '';
|
||||
const TURN_TIMEOUT_MS = Number(process.env.AGENTMAIL_TURN_TIMEOUT_MS || 60_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);
|
||||
|
||||
// ─── 进程内状态 ───
|
||||
//
|
||||
// 全部只在内存,重启即丢 —— 这是契约第六节列明的已知取舍。
|
||||
// 要持久化的话该落在 pi 的会话元数据里,而不是桥自己的文件。
|
||||
|
||||
const sessions = new Map(); // agentmail session_id -> { session, sessionManager, cwd }
|
||||
const reverseMap = new Map(); // pi session id -> agentmail session_id
|
||||
const mailDriven = new Set(); // pi session id
|
||||
const mailContexts = new Map(); // agentmail session_id -> { replyTo, subject, mailID }
|
||||
const relayedSummaries = new Map(); // pi session id -> 已转发过的 relay_key
|
||||
const syncedNames = new Map(); // pi session id -> 上次提交给 Gateway 的名字
|
||||
const pendingPermissions = new Map(); // relay_key -> { resolve, piSessionId }
|
||||
const deliveredMails = new Set(); // 已投过的 mail_id(SSE 与补拉共用,B-7.3)
|
||||
|
||||
let allowedModels = [];
|
||||
let modelRuntime = null;
|
||||
let client = 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-8)───
|
||||
|
||||
/**
|
||||
* 内联 pi 扩展:把 pi 拦下的危险工具调用转成一封邮件问人。
|
||||
*
|
||||
* 这是 `I-1` 最直接的体现 —— 被平台真正拦下的那一次才是事实,
|
||||
* 不依赖模型「记得」调 request_permission(它会忘,也会在不需要时乱调)。
|
||||
*
|
||||
* pi 的 `tool_call` 钩子**可以 await**(C-9 实测成立:处理器里 await 300ms
|
||||
* 再返回 {block:true},pi 会等),所以这里能真的等人做决定,
|
||||
* 不必走「先拒一次再重试」的退化路径。
|
||||
*
|
||||
* @param {string} piSessionIdRef 用一个 getter 拿会话 id:扩展工厂在
|
||||
* createAgentSession **内部**被调用,那时 session 对象还没返回给桥。
|
||||
*/
|
||||
function permissionExtension(getMailContext) {
|
||||
// pi 默认放行内建工具;桥只拦真正有副作用的那几个。
|
||||
// read/grep/ls 之类不拦:每一步都问人会让 Agent 什么也做不成,
|
||||
// 而人也会很快开始无脑点同意(那比不问更危险)。
|
||||
const GUARDED = new Set(['bash', 'write', 'edit']);
|
||||
|
||||
return (pi) => {
|
||||
pi.on('tool_call', async (event, ctx) => {
|
||||
if (!GUARDED.has(event.toolName)) return;
|
||||
|
||||
const piSessionId = ctx?.sessionManager?.getSessionId?.() || '';
|
||||
const mailSessionId = reverseMap.get(piSessionId);
|
||||
// 不是邮件驱动的会话 → 让位给 pi 自己的本地 UI(B-8.2)。
|
||||
// 占着钩子不放会让人在 TUI 里干活时每一步都卡住等邮件。
|
||||
if (!mailSessionId) return;
|
||||
|
||||
// relay_key 用 pi 给的 toolCallId(B-8.1):服务端会随决策事件回传它,
|
||||
// 桥重启丢了 pendingPermissions 也能对上(B-4.2)。自造随机 id 做不到。
|
||||
const relayKey = `${piSessionId}:${event.toolCallId}`;
|
||||
const ctxInfo = getMailContext(mailSessionId);
|
||||
|
||||
try {
|
||||
await client.post('/permission/request', {
|
||||
question: `是否允许执行 ${event.toolName}?`,
|
||||
options: ['同意', '一直同意', '拒绝'],
|
||||
context: describeToolCall(event),
|
||||
session_id: mailSessionId,
|
||||
to: ctxInfo?.replyTo || '',
|
||||
relay_key: relayKey,
|
||||
});
|
||||
} catch (e) {
|
||||
// 转发失败 → 让位给 pi 本地 UI(B-8.2)。返回 undefined 表示
|
||||
// 「这个钩子不表态」,pi 会走它自己的批准流程。
|
||||
log(`权限转发失败,让位给本地决策: ${describeError(e)}`);
|
||||
return;
|
||||
}
|
||||
|
||||
log(`权限询问已发出(${event.toolName},key=${relayKey}),等待决策…`);
|
||||
const decision = await new Promise((resolve) => {
|
||||
pendingPermissions.set(relayKey, { resolve, piSessionId });
|
||||
});
|
||||
|
||||
// fail closed(B-9.2 / N-9):只有明确的同意才放行。
|
||||
// 关停时 shutdown() 会用 'shutdown' 唤醒所有等待者,落到这里的 else。
|
||||
if (/^(同意|一直同意|allow|approve|always|yes)/i.test(decision)) {
|
||||
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);
|
||||
}
|
||||
|
||||
// ─── 会话解析(B-3)───
|
||||
|
||||
/**
|
||||
* 没有可用 `to_workspace` 时的兜底目录。
|
||||
*
|
||||
* 与 DSH 的 `mailSessionFallback` 同构,但目录名是 `.pi`:那个函数在
|
||||
* lib/ 下(三平台逐字节相同),写死了 `.dsh`,不能为 pi 改。
|
||||
* 让 pi 的会话落进 `~/.dsh/` 会让人以为是 DSH 在干活。
|
||||
*/
|
||||
function piMailFallback(sessionKey) {
|
||||
return join(homedir(), '.pi', 'mail-sessions', String(sessionKey || 'default'));
|
||||
}
|
||||
|
||||
/**
|
||||
* 找到(或建立)这封邮件该落进的 pi 会话。
|
||||
*
|
||||
* Gateway 已经按三维地址的 session 位做完了「复用默认 / 新建 / 具名必须存在」
|
||||
* 的判定,推来的 session_id 就是判定结果 —— 桥只负责忠实映射,
|
||||
* 不自己决定开不开新会话(N-8:404 后自动改用 .new 是禁止的)。
|
||||
*/
|
||||
async function resolveSession(data, mailTools) {
|
||||
const mailSessionID = data.session_id;
|
||||
const bound = mailSessionID ? sessions.get(mailSessionID) : undefined;
|
||||
if (bound) return { ...bound, reused: true };
|
||||
|
||||
// cwd 取寻址里的 path 位(B-3.1)。校验走共用模块:目录不存在时**不创建**
|
||||
// (N-2:笔误会在磁盘上落下真目录,而 Agent 在里面一无所获),拒绝相对路径(N-3)。
|
||||
//
|
||||
// 兜底用 `~/.pi/mail-sessions/<会话>` 而不是共用模块里的 mailSessionFallback ——
|
||||
// 后者写死了 `.dsh` 目录名(那是 DSH 的家),pi 的会话落进去会让人以为
|
||||
// DSH 在干活。lib/ 里的函数三平台逐字节相同,不能为 pi 改它。
|
||||
const { cwd, grouped } = resolveWorkspaceCwd(data.to_workspace, piMailFallback(mailSessionID));
|
||||
if (!grouped && data.to_workspace) {
|
||||
log(`工作目录 ${data.to_workspace} 不可用,回退到 ${cwd}`);
|
||||
}
|
||||
ensureCwd(cwd, grouped);
|
||||
|
||||
const opened = await openSession({
|
||||
cwd,
|
||||
modelRuntime,
|
||||
customTools: mailTools,
|
||||
extension: permissionExtension((id) => mailContexts.get(id)),
|
||||
});
|
||||
for (const d of opened.diagnostics) {
|
||||
log(`扩展诊断: ${d?.message ?? JSON.stringify(d)}`);
|
||||
}
|
||||
|
||||
const piSessionId = opened.session.sessionId;
|
||||
const entry = { session: opened.session, sessionManager: opened.sessionManager, cwd };
|
||||
|
||||
if (mailSessionID) {
|
||||
sessions.set(mailSessionID, entry);
|
||||
reverseMap.set(piSessionId, mailSessionID);
|
||||
mailDriven.add(piSessionId);
|
||||
}
|
||||
|
||||
// 一轮结束就转发总结(B-5)。挂 agent_end 而不是 message_end:
|
||||
// 后者在流式生成中反复触发,转出去的是半截话。
|
||||
// subscribe 收的是一个普通函数(AgentSessionEventListener),不是 {onEvent}。
|
||||
opened.session.subscribe((event) => {
|
||||
if (event?.type === 'agent_end') {
|
||||
// willRetry 为真表示 pi 自己要重试(auto_retry),这一轮还没定论 —— 不转。
|
||||
if (event.willRetry) return;
|
||||
relaySummary(piSessionId).catch((e) => log(`自动转发失败: ${describeError(e)}`));
|
||||
}
|
||||
// pi 侧改名(pi-web 生成标题、人在 TUI 里 /name)→ 同步给 Gateway
|
||||
if (event?.type === 'session_info_changed') {
|
||||
syncNaming(piSessionId, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`));
|
||||
}
|
||||
});
|
||||
|
||||
log(`新建 pi 会话 ${piSessionId}(cwd=${cwd})`);
|
||||
return { ...entry, reused: false };
|
||||
}
|
||||
|
||||
// ─── 命名一致(C-11 / W-7)───
|
||||
|
||||
/**
|
||||
* pi 的名字 → Gateway → 定稿别名回写进 pi。
|
||||
*
|
||||
* 完整推理见 src/naming.mjs 顶部。这里只是把那套决策接上 I/O。
|
||||
*/
|
||||
async function syncNaming(piSessionId, platformName) {
|
||||
const mailSessionID = reverseMap.get(piSessionId);
|
||||
if (!mailSessionID) return; // 不是邮件驱动的会话,不碰
|
||||
|
||||
const plan = planNamingSync({
|
||||
platformName,
|
||||
mailSubject: mailContexts.get(mailSessionID)?.subject,
|
||||
lastSynced: syncedNames.get(piSessionId),
|
||||
});
|
||||
if (plan.skip) return;
|
||||
|
||||
// 先记下指纹再发请求:响应回来时 setSessionName 会再次触发
|
||||
// session_info_changed,这一步是防自激循环的关键。
|
||||
syncedNames.set(piSessionId, plan.signature);
|
||||
|
||||
const res = await client.post(`/sessions/${mailSessionID}/sync`, {
|
||||
alias: plan.alias,
|
||||
title: plan.title,
|
||||
});
|
||||
|
||||
const entry = sessions.get(mailSessionID);
|
||||
const back = planWriteBack({
|
||||
finalAlias: res?.alias,
|
||||
currentPiName: entry?.session?.sessionName,
|
||||
});
|
||||
log(`命名同步 ${piSessionId}: alias=${res?.alias || '(未变)'} 来源=${plan.source}`);
|
||||
|
||||
if (back.write && entry?.session) {
|
||||
// 顺序要紧:先更新指纹,再改名。
|
||||
//
|
||||
// setSessionName **同步**触发 session_info_changed(实测),于是本函数会在
|
||||
// 这一行里被重入。指纹在改名之后才更新的话,重入那次看到的还是旧指纹,
|
||||
// 于是又打一次 sync —— 每条会话两次请求,内容完全相同。
|
||||
//
|
||||
// 记的是「把定稿别名当作平台名字」会算出的指纹:重入那次的 platformName
|
||||
// 正是 back.name,来源判定成 platform,算出来的就是这个值。
|
||||
syncedNames.set(piSessionId, `platform:${back.name}|${back.name}`);
|
||||
// 只用 setSessionName(走 pi 自己的写入路径)。绝不自己拼路径写会话文件:
|
||||
// 首条 assistant 消息落盘前文件还不存在,pi 首次落盘用 openSync(file,"wx"),
|
||||
// 抢先创建会让它抛 EEXIST(实测)。
|
||||
entry.session.setSessionName(back.name);
|
||||
log(`别名回写 pi:${back.name}(${back.reason})`);
|
||||
}
|
||||
}
|
||||
|
||||
// ─── 自动转发(B-5)───
|
||||
|
||||
async function relaySummary(piSessionId) {
|
||||
const mailSessionID = reverseMap.get(piSessionId);
|
||||
if (!mailSessionID) return;
|
||||
// 只对邮件驱动的会话转发(B-5.5):人在 pi 里正常干活时不该往邮箱灌总结
|
||||
if (!mailDriven.has(piSessionId)) return;
|
||||
|
||||
const entry = sessions.get(mailSessionID);
|
||||
if (!entry) return;
|
||||
|
||||
// 一轮结束是命名的自然时机(C-11 / D-5)。
|
||||
//
|
||||
// 这一步不能只挂在 session_info_changed 上:桥用 SDK 起的会话**永远不会**
|
||||
// 触发那个事件 —— pi 的标题生成器在 pi-web 里,不在内核里,SDK 路径上没有它。
|
||||
// 只等事件的话别名永远是空的,于是 `name@path.<别名>` 续谈无从下手
|
||||
// (实测过:第一封邮件跑通了,sessions.session_alias 仍是空串)。
|
||||
//
|
||||
// 放在转发**之前**:回信里会带上会话别名,收件人看到的第一封回信就能用它续谈。
|
||||
await syncNaming(piSessionId, entry.session.sessionName)
|
||||
.catch((e) => log(`命名同步失败: ${describeError(e)}`));
|
||||
|
||||
// 只取 type==='text' 的块(B-5.1 / N-6):thinking 是思考过程,不是结论
|
||||
const text = lastAssistantText(entry.session.messages);
|
||||
if (!text) return; // 空文本不发空邮件(B-5.4)
|
||||
|
||||
const ctx = mailContexts.get(mailSessionID);
|
||||
if (!ctx?.replyTo) return; // 不知道回给谁
|
||||
|
||||
// 幂等键用 pi 的会话 id + 会话树叶子 id:两者都落盘,重启重放也是同一个键。
|
||||
const relayKey = relayKeyFor(piSessionId, entry.sessionManager.getLeafId?.());
|
||||
if (relayedSummaries.get(piSessionId) === relayKey) return;
|
||||
|
||||
// 模型这一轮已亲手回过这条线索 → 让位(B-5.3)。
|
||||
// 否则收件箱里是两封说同一件事的邮件(生产实测过)。
|
||||
if (shouldSkipAutoRelay(explicitSends.get(piSessionId), ctx.replyTo, ctx.mailID)) {
|
||||
explicitSends.delete(piSessionId);
|
||||
relayedSummaries.set(piSessionId, relayKey);
|
||||
log(`本轮模型已主动回信 ${ctx.replyTo},跳过自动转发`);
|
||||
return;
|
||||
}
|
||||
|
||||
await client.post('/mail/send', {
|
||||
to: ctx.replyTo,
|
||||
subject: replySubject(ctx.subject),
|
||||
body: text,
|
||||
reply_to: ctx.mailID || '',
|
||||
// relay + relay_key 走免配额通道(I-2):模型已经把话说完了,
|
||||
// 桥只是把它搬到邮件里。对搬运收费会让配额用尽时 Agent 连交代都做不了。
|
||||
relay: 'summary',
|
||||
relay_key: relayKey,
|
||||
});
|
||||
relayedSummaries.set(piSessionId, relayKey);
|
||||
explicitSends.delete(piSessionId); // 一轮结束,窗口关闭
|
||||
log(`已转发本轮总结给 ${ctx.replyTo}(${text.length} 字)`);
|
||||
}
|
||||
|
||||
// ─── 投递(B-3 / B-6)───
|
||||
|
||||
async function deliverMail(data, kind, mailTools) {
|
||||
const { session, reused } = await resolveSession(data, mailTools);
|
||||
const piSessionId = session.sessionId;
|
||||
|
||||
// 新一轮开始:清掉上一轮「模型主动发过信」的记录。不清的话,
|
||||
// 上一轮亲手回过信会永久压掉这个会话之后所有的自动转发。
|
||||
explicitSends.delete(piSessionId);
|
||||
|
||||
if (kind === 'mail' && data.session_id) {
|
||||
// 一个会话里可能来过多封信,只留最近那封 —— 回信要落回最新的线索
|
||||
mailContexts.set(data.session_id, {
|
||||
replyTo: data.from_name || '',
|
||||
subject: data.subject || '',
|
||||
mailID: data.mail_id || '',
|
||||
});
|
||||
}
|
||||
|
||||
const prompt = buildMailPrompt({ agentName: AGENT_NAME, data, kind, reused });
|
||||
|
||||
// 续谈:会话已经存在,模型也已经定了(pi 的模型在 createAgentSession 时绑定),
|
||||
// 所以这一支不做模型降级。runTurn 内部按 isStreaming 分流:
|
||||
// 空闲就直接起一轮,正在跑就排到当轮之后(不打断上一封邮件的工作)。
|
||||
if (reused) {
|
||||
const outcome = await runTurn(session, prompt, TURN_TIMEOUT_MS);
|
||||
log(`续谈 ${piSessionId}(mail ${data.mail_id}${outcome.queued ? ',已排队' : ''})`);
|
||||
// 续谈失败不换模型重试(换模型要换会话,会丢掉整条上下文 ——
|
||||
// 而上下文正是发件人指定这条会话的原因),但要让失败可见。
|
||||
if (!outcome.ok) throw new Error(`续谈失败: ${outcome.error}`);
|
||||
return;
|
||||
}
|
||||
|
||||
// 按管理员划定的范围逐个尝试(D-3)。
|
||||
// 关键点:`prompt()` resolve **不代表模型跑成功了** —— 无凭证的 provider
|
||||
// 会让它 reject(实测 `No API key found for amazon-bedrock.`),
|
||||
// 而上游报错走 stopReason==='error'。判定交给 classifyTurnOutcome。
|
||||
const attempts = modelAttemptOrder(allowedModels, {
|
||||
provider: REPLY_PROVIDER,
|
||||
model: REPLY_MODEL,
|
||||
});
|
||||
const failures = [];
|
||||
|
||||
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 消息,丢掉不损失内容。
|
||||
const cwd = sessions.get(data.session_id)?.cwd;
|
||||
const current = sessions.get(data.session_id)?.session;
|
||||
current?.dispose?.();
|
||||
const retried = await openSession({
|
||||
cwd,
|
||||
modelRuntime,
|
||||
model,
|
||||
customTools: mailTools,
|
||||
extension: permissionExtension((id) => mailContexts.get(id)),
|
||||
});
|
||||
rebind(data.session_id, current?.sessionId ?? piSessionId, retried, cwd);
|
||||
const outcome = await runTurn(retried.session, prompt, TURN_TIMEOUT_MS);
|
||||
if (outcome.ok) {
|
||||
if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`);
|
||||
return;
|
||||
}
|
||||
failures.push({ ...route, error: outcome.error });
|
||||
log(`模型 ${label} 失败: ${outcome.error}`);
|
||||
continue;
|
||||
}
|
||||
|
||||
const outcome = await runTurn(session, prompt, TURN_TIMEOUT_MS);
|
||||
if (outcome.ok) {
|
||||
if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`);
|
||||
return;
|
||||
}
|
||||
failures.push({ 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(' | ')}`);
|
||||
}
|
||||
|
||||
/** 换模型重开会话后,把三张映射表指向新会话。 */
|
||||
function rebind(mailSessionID, oldPiId, opened, cwd) {
|
||||
reverseMap.delete(oldPiId);
|
||||
mailDriven.delete(oldPiId);
|
||||
const piSessionId = opened.session.sessionId;
|
||||
// cwd 由调用方传:AgentSession 上没有 cwd getter(只有 sessionId /
|
||||
// sessionFile / sessionName),从 sessionManager.getCwd() 也行,
|
||||
// 但这里本来就有那个值,多绕一层没有意义。
|
||||
const entry = { session: opened.session, sessionManager: opened.sessionManager, cwd };
|
||||
if (mailSessionID) {
|
||||
sessions.set(mailSessionID, entry);
|
||||
reverseMap.set(piSessionId, mailSessionID);
|
||||
mailDriven.add(piSessionId);
|
||||
}
|
||||
opened.session.subscribe((event) => {
|
||||
if (event?.type === 'agent_end' && !event.willRetry) {
|
||||
relaySummary(piSessionId).catch((e) => log(`自动转发失败: ${describeError(e)}`));
|
||||
}
|
||||
if (event?.type === 'session_info_changed') {
|
||||
syncNaming(piSessionId, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// ─── 权限决策回来(B-4)───
|
||||
|
||||
async function handlePermissionDecision(data, mailTools) {
|
||||
const relayKey = data.relay_key || '';
|
||||
const pending = relayKey ? pendingPermissions.get(relayKey) : undefined;
|
||||
|
||||
if (pending) {
|
||||
pendingPermissions.delete(relayKey);
|
||||
pending.resolve(String(data.decision || '拒绝'));
|
||||
log(`权限 ${relayKey} 决策 ${data.decision}(决策人 ${data.decided_by || '?'})`);
|
||||
return;
|
||||
}
|
||||
|
||||
// 找不到挂起项(桥重启丢了内存映射)→ 退化为把决策当一封通知投进原会话(B-4.2)。
|
||||
// 此时 pi 侧那次工具调用早已随进程消失,但人刚刚点了「同意」——
|
||||
// 什么都不做的话人以为自己批准了、Agent 却毫无反应。
|
||||
if (!data.session_id || !sessions.has(data.session_id)) {
|
||||
// **不得凭空新开会话**(B-4.3)
|
||||
log(`权限决策 ${relayKey} 无对应会话,忽略`);
|
||||
return;
|
||||
}
|
||||
log(`权限 ${relayKey} 无挂起项,退化为通知投递`);
|
||||
await deliverMail(data, 'permission', mailTools);
|
||||
}
|
||||
|
||||
// ─── 心跳(B-2)───
|
||||
|
||||
async function reportSessions() {
|
||||
try {
|
||||
const { SessionManager } = await import('@earendil-works/pi-coding-agent');
|
||||
// 不传参数:`listAll(dir)` 把字符串当**自定义会话目录**,传 getAgentDir()
|
||||
// 会去 ~/.pi/agent 下直接找 .jsonl(那里没有),得到空列表。
|
||||
// 不传时它用默认的 ~/.pi/agent/sessions,逐个 cwd 子目录扫。
|
||||
//
|
||||
// 用 listAll 而不是 list(cwd):桥的进程 cwd 与会话 cwd 无关,
|
||||
// 按前者过滤会漏掉所有真正在干活的会话。
|
||||
const all = await SessionManager.listAll();
|
||||
return snapshotPiSessions(all, (id) => mailDriven.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;
|
||||
}
|
||||
}
|
||||
|
||||
async function catchUp(pending, mailTools) {
|
||||
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} 封未读)`);
|
||||
// 串行(B-7.2):每封都要起一轮模型,并发放出去等于对上游打 N 个并发请求
|
||||
for (const ev of tasks) {
|
||||
if (deliveredMails.has(ev.mail_id)) continue; // 逐封再查(B-7.6)
|
||||
deliveredMails.add(ev.mail_id);
|
||||
try {
|
||||
await deliverMail(ev, 'mail', mailTools);
|
||||
} catch (e) {
|
||||
log(`补投 ${ev.mail_id} 失败: ${describeError(e)}`);
|
||||
}
|
||||
}
|
||||
} 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 建一次全进程共用:它要读 auth.json / models.json 并做
|
||||
// 可用性探测,每条会话建一个既慢又会重复打 provider 的探测请求。
|
||||
//
|
||||
// allowModelNetwork 保持默认的 false:桥启动时不去网上拉模型目录。
|
||||
// 拉了也没用 —— 上报给 Gateway 的是 getAvailable()(有凭证、真能调起来的),
|
||||
// 而那取决于本机 auth.json,不取决于目录里有多少条。开着只会让
|
||||
// 启动多等一个网络往返,而且断网时启动路径上多一个可失败点。
|
||||
modelRuntime = await ModelRuntime.create();
|
||||
const runtimeErr = modelRuntime.getError?.();
|
||||
if (runtimeErr) log(`模型运行时告警: ${runtimeErr}`);
|
||||
|
||||
const mailTools = createMailTools({ client, log, agentName: AGENT_NAME });
|
||||
|
||||
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, mailTools);
|
||||
}
|
||||
} catch {
|
||||
// B-2.1:心跳失败不重试不报错。真连不上时 Gateway 会把它判成离线,
|
||||
// 那才是可见的信号;桥自己打一串错误日志只会淹掉真正的问题。
|
||||
}
|
||||
};
|
||||
await beat(); // B-1.3:不等第一个 30 秒周期
|
||||
heartbeatTimer = setInterval(beat, 30_000); // B-1.5
|
||||
|
||||
client.startSSE((type, data) => { // B-1.4:首次不带 Last-Event-ID
|
||||
if (type === 'permission_decision') {
|
||||
handlePermissionDecision(data, mailTools).catch((e) =>
|
||||
log(`权限决策处理失败: ${describeError(e)}`));
|
||||
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);
|
||||
deliverMail(data, 'mail', mailTools).catch((e) => log(`投递 ${id} 失败: ${describeError(e)}`));
|
||||
}, log);
|
||||
|
||||
for (const sig of ['SIGINT', 'SIGTERM']) process.on(sig, () => shutdown(sig));
|
||||
}
|
||||
|
||||
function shutdown(reason) {
|
||||
if (shuttingDown) return;
|
||||
shuttingDown = true;
|
||||
log(`收到 ${reason},关停中…`);
|
||||
|
||||
if (heartbeatTimer) clearInterval(heartbeatTimer); // B-9.1
|
||||
client?.stopSSE();
|
||||
|
||||
// B-9.2 / N-9:所有未决权限询问 fail closed。
|
||||
// 不唤醒的话 pi 侧那些 await 永不返回,整条会话挂死;
|
||||
// 而默认放行一个没人批准的危险操作,比让它失败严重得多。
|
||||
for (const [key, p] of pendingPermissions) {
|
||||
log(`未决权限 ${key} fail closed`);
|
||||
p.resolve('shutdown');
|
||||
}
|
||||
pendingPermissions.clear();
|
||||
|
||||
for (const { session } of sessions.values()) {
|
||||
try { session.dispose?.(); } catch { /* 关停期的报错没有价值 */ }
|
||||
}
|
||||
releaseLock();
|
||||
// B-9.3:不发「插件下线」通知邮件
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
log(`启动失败: ${describeError(e)}`);
|
||||
releaseLock();
|
||||
process.exit(1);
|
||||
});
|
||||
Reference in New Issue
Block a user