按用户裁定「yolo_own_tools」实现:平台让开(--mode yolo),它自带的一切
「能动机器」的工具被 --disallowed-tools 拿掉,执行类动作改由我们自己的
run_command / write_file 承担,而门禁就在这两个工具里 —— 逐次向发件人请示。
## 为什么必须走这条路(实测,不是推断)
MCP 工具的 needsApproval 在产物里**硬编码为 true**(与 annotations 无关),
而 build/edit 档的判定最后一条是「需要审批 → ask」;headless 没有审批客户端
可问 ⇒ **每个 MCP 工具都被拒**(连 read_inbox 都调不动)。
我们本想让平台把询问转给钩子,但 PermissionRequest 在本版本(3.10.2 / CLI 0.16.5)
**不可靠**:有时压根不注册,触发时也无条件在 ~5ms 内失败、命令从未被 spawn
(用「钩子写 marker 文件」的副作用验证)。
于是选择只剩两个:「平台问、但问不到人 → 全拒」与「平台不问、我们自己问」。
后者才既可用又可审计。代价(平台不再提供第二道防线)写进了 README 的残余风险。
## 新增
- `lib/approval.mjs`:授权往返的唯一实现(钩子与工具共用,否则必然漂移)。
三条不可动摇的规矩:只有明确同意才放行(判据是共用库的前缀白名单,
不是「不等于拒绝」);永久失败(409/4xx)当场拒绝并把服务端建议带给模型;
暂时失败看有没有本地界面 —— 判据用**调用方传的 sessionId**(单一事实来源,
不再另读环境变量)。自己开 SSE 等决定,先建连再发请求。
- `lib/action-tools.mjs`:`run_command` / `write_file`。输出上限、超时上限、
默认 cwd=工作区;拒绝时**抛错**(MCP 层转 isError)而不是返回「已处理」——
opencode 上「工具失败但报成功」导致模型连试 6 次后放弃整个任务的教训。
平台保护目录(网关数据库/插件代码/服务单元/密钥目录)**无论谁批准都不写**,
且判定在门禁之前(不消耗人的注意力)——防的是自我强化:邮件驱动的 Agent
可能被来信诱导去改自己的插件代码,改完下一轮就换了一套规则。
- `REVIEWED_DENYLIST`(32 项):逐条按「不拿掉会怎样」分类。名单来自 CLI 产物里
模型可见工具名的**权威注册表**(aIn 那个 28 项数组)+ 另一份更宽的候选集并集,
**不采信模型自述**(基线里它用某个没点名的方式真的创建了文件)。
最容易被漏掉的是 `js` / `mcp__node_repl__js`:它挂在 MCP 上、
产物里自述「can run arbitrary JavaScript with full Node privileges, like Bash」。
- 提示词的能力说明(分档):告诉模型自带工具被禁、动手要用哪两个工具、
会被请示;并明确「被拒是业务结果,不要重试、不要绕道」。
## 修掉三个真缺陷(都是实测撞出来的)
1. **幂等键按「会话+工具」取 → 同会话第二次调用被静默吞掉**。
网关对重复 relay_key 返回 **HTTP 200** `{status:"duplicate_relay"}` 并提前返回:
不建请求、不发邮件、**永远不会有人来决策**。于是工具干等 → 被 MCP 调用超时
砍掉 → 模型回报「30 秒内未获批准」。从状态码到措辞全看不出问题,归因还完全
错了(像是人没理它)。改为**按调用唯一**(保留会话/工具前缀便于反查),
并把 duplicate_relay 当成可读的拒绝(fail fast,不再干等)。
2. **授权窗口被 MCP 调用超时截断**。ZCode 对 MCP 工具调用有超时(默认量级 30 秒),
而门禁要等人。已在插件清单声明 `mcpServers.agentmail.timeoutMs=600000`
(实测生效:40 秒的命令没被砍,墙钟 50 秒通过),并让门禁**自己**把等待夹到
timeoutMs - 余量之下(`resolveWaitMs`)——被客户端杀掉时连理由都发不出去,
所以必须由我们自己先 settle。
3. **`--allowed-tools` 在 help 里写着但解析器不认**(`Unknown option`)。
留着会拼出一条永远跑不起来的命令行,现在 `buildRunArgs` 直接抛错并指出
替代方案。我在这里误判过一次:先看到「文件没创建」就以为白名单生效,
其实进程只是没退到 usage。判据缺了「进程真的执行了」这一环。
## 自报改成如实
detectModeEnforcement 以前拿「钩子已注册」当 native 的凭据 —— yolo 下钩子
根本不会触发,那等于替一个不存在的能力背书。现在先看**我们那条链**是否就绪
(yolo + 禁用清单里真的有 Bash/js),就绪才报 native,并在理由里点明谁在把关
(实测输出:「执行类动作只能经我们自己的门禁…平台自带危险工具已禁用 32 项」)。
## 验证
- 单测 376/376(新增 47 条)。重点在反向对照:一句「拒绝/deny/空串/平台自己的
shutdown 哨兵都不放行」之外,还验了「别人的决策不能拿来用(relay_key 配对)」、
「超时必须真的拒绝」、「同一会话两次调用必须用不同的幂等键」、
「重复请求要当场拒绝而不是干等」;执行工具的每条拒绝场景都配一个**文件系统断言**
(「抛错了」不等于「副作用没发生」),保护目录还验了 `..`/`./` 绕不过去。
- 真模型端到端(`/root/e2e-zcode-gate/run.py`,13/13):
批 → 命令真执行(文件内容=标记);拒 → 命令真没执行(文件不存在)
且回信把成因说成「人拒绝」而**不是**「超时」;同会话第三次调用仍能产生新请求
并在获批后执行。判据本身也修了两处(授权请求邮件里带标记会被误当成回信;
备注在通过项旁边显示会误导)。
- 部署:`deploy/redeploy-plugin.sh zcode` 快照切换 + 握手自检;
驱动单元改为跑快照(生产不跑仓库工作区),env 与清单超时的关系写进注释。
- 顺手清掉一个遗留驱动进程(跑的是仓库路径的旧代码、连着网关 SSE、会抢邮件)。
## 判据纪律(本轮又踩到、已写进代码注释)
「文件没被创建」不能区分「被拦住了」与「进程根本没跑」;
「未获批准」不能区分「人拒绝」与「窗口被截断」;
「工具报错」不能区分「命令失败」与「工具坏了」。
每一处都改成了验到**具体成因**。
474 lines
19 KiB
JavaScript
474 lines
19 KiB
JavaScript
#!/usr/bin/env node
|
||
/**
|
||
* ZCode 的 AgentMail 驱动:**收到来信 → 起一轮 ZCode → 把结论回信**。
|
||
*
|
||
* 这是让 ZCode 成为一等 Agent 的那一半(另一半是插件:MCP 工具面 + 授权钩子)。
|
||
*
|
||
* # 与另三个桥的关系
|
||
*
|
||
* 结构对齐 pi / dsh / opencode 三桥:SSE 订阅 → 去重 → 解析工作目录与档位 →
|
||
* 跑一轮 → 按策略回信 → 心跳。可复用的部分一律走 `lib/`(逐字节同源):
|
||
* 事件补投、工作目录解析、回信策略、去重判据、SSE 帧解析。
|
||
*
|
||
* 差别只在「怎么跑一轮」:ZCode 用 **headless CLI**
|
||
* (`--prompt … --output-format stream-json`),不是 SDK。
|
||
*
|
||
* # 三个必须记住的约束
|
||
*
|
||
* 1. **`--mode` 必传**。`--prompt` 的默认 mode 是 `yolo`,而 yolo 会绕过全部
|
||
* 权限询问 —— 授权钩子根本不会触发,授权系统会**静默消失**(不报错,
|
||
* 只是没有任何询问)。档位映射见 `src/turn-mode.mjs`。
|
||
* 2. **失败必须回信**。邮件驱动的会话没有本地界面,一轮跑不起来而什么都不发,
|
||
* 发件人只会觉得「信发出去了,然后再无音讯」。
|
||
* 3. **模型自己发过信就不再自动转发**。工具跑在 ZCode 派生的 MCP 服务器进程里,
|
||
* 与驱动不是同一个进程,所以经 `lib/explicit-sends.mjs` 落盘对齐。
|
||
*
|
||
* # 串行
|
||
*
|
||
* 一轮一次。ZCode 的会话与工作目录是重资源,同一目录并发跑两轮会互相踩;
|
||
* 代价是一封长信会挡住后面的信 —— 这是显式取舍,不是遗漏(见 README 的已知缺口)。
|
||
*/
|
||
|
||
import { basename, join } from 'node:path';
|
||
import { homedir } from 'node:os';
|
||
import { readFileSync } from 'node:fs';
|
||
import { fileURLToPath } from 'node:url';
|
||
import { GatewayClient, GatewayError } from '../lib/gateway.mjs';
|
||
import { createSSEClient } from '../lib/sse-client.js';
|
||
import { BoundedSet, BoundedMap, MAX_TRACKED_MAILS, MAX_TRACKED_SESSIONS } from '../lib/bounded.js';
|
||
import { selectCatchup } from '../lib/catchup.js';
|
||
import { autoRelayDecision } from '../lib/relay-policy.js';
|
||
import { shouldSkipAutoRelay } from '../lib/relay-dedup.js';
|
||
import { resolveWorkspaceCwd, ensureCwd } from '../lib/workspace.js';
|
||
import { normalizeMode } from '../lib/permission-mode.js';
|
||
import { clampRelayKey } from '../lib/relay-key.js';
|
||
import { explicitSendsFile, readExplicitSends } from '../lib/explicit-sends.mjs';
|
||
import { zcodeModeForTier, denylistForTier, modeReachesPermissionHook, describeTier } from './turn-mode.mjs';
|
||
import { buildMailPrompt, replySubject, renderTurnFailure } from './prompt.mjs';
|
||
import { runTurn, DEFAULT_CLI } from './zcode-run.mjs';
|
||
import { isMainModule } from '../lib/is-main.mjs';
|
||
|
||
const log = (...parts) => console.error('[zcode-mail-bridge]', ...parts);
|
||
|
||
const CONFIG = {
|
||
gatewayURL: process.env.AGENTMAIL_GATEWAY_URL || 'http://127.0.0.1:8180',
|
||
agentName: process.env.AGENTMAIL_AGENT_NAME || 'zcode',
|
||
turnTimeoutMs: Number(process.env.AGENTMAIL_TURN_TIMEOUT_MS || 20 * 60 * 1000),
|
||
maxTurns: Number(process.env.AGENTMAIL_MAX_TURNS || 0) || undefined,
|
||
workspaceRoot: process.env.AGENTMAIL_WORKSPACE_ROOT || '',
|
||
cliPath: process.env.AGENTMAIL_ZCODE_CLI || DEFAULT_CLI
|
||
};
|
||
|
||
/**
|
||
* 没有 `to_workspace` 时的兜底目录。
|
||
*
|
||
* **不能用共用的 `mailSessionFallback`** —— 那个函数的目录名写的是 `~/.dsh`
|
||
* (它注释里也写明是「没有天然兜底的平台(DSH)用这个」)。各平台的会话存储
|
||
* 各不相同,把 ZCode 的会话塞进 `~/.dsh` 下会造成两个平台的会话目录互相污染。
|
||
*
|
||
* @param {string} sessionKey
|
||
*/
|
||
export function zcodeSessionFallback(sessionKey, rootOverride) {
|
||
const root = rootOverride || CONFIG.workspaceRoot;
|
||
if (root) return join(root, String(sessionKey || 'default'));
|
||
return join(homedir(), '.zcode', 'mail-sessions', String(sessionKey || 'default'));
|
||
}
|
||
|
||
/**
|
||
* 自报给网关的强制力。
|
||
*
|
||
* **只声明得出来的事**。而且要分开两件事,它们以前被当成了同一件:
|
||
*
|
||
* 1. **平台自己的权限引擎**(PermissionRequest 钩子)—— 只在 `--mode build/edit`
|
||
* 下才会问。现在档位映射把 workspace 映到 `yolo`(平台恒 allow、钩子根本不会触发),
|
||
* 所以**拿「钩子已注册」当强制力证据是错的**。桌面(有界面)时它仍然有用,
|
||
* 但 headless 驱动这条路上不是它。
|
||
* 2. **我们自己的门禁**(`lib/action-tools.mjs`)—— `yolo` + `--disallowed-tools`
|
||
* 之后,能动机器的只剩我们两个工具,而它们每次都要过 `requestApproval`。
|
||
* 这个才是这条路上真正拦得住东西的那一层。
|
||
*
|
||
* 于是判据改成:先看我们自己那条链能不能真的拦住(模块在不在 + 禁用清单够不够),
|
||
* 再说平台那一层是否还参与。报 native 的含义是「该档位真的被强制,而非口头约定」。
|
||
*
|
||
* 查的是驱动自己的文件(它住在插件里),不需要额外配置项。
|
||
*/
|
||
export function detectModeEnforcement({ hooksFile, env = process.env } = {}) {
|
||
const file = hooksFile || fileURLToPath(new URL('../hooks/hooks.json', import.meta.url));
|
||
|
||
// 我们自己那条链:必须是 yolo(平台不问)+ 一张够长的禁用清单。
|
||
const mode = zcodeModeForTier('workspace', env);
|
||
const denied = denylistForTier('workspace', env);
|
||
const gateReady = mode === 'yolo' && denied.includes('Bash') && denied.includes('js');
|
||
|
||
// 平台那条链:钩子清单里真的注册了 PermissionRequest(桌面模式下才用得上)。
|
||
let hookRegistered = false;
|
||
let hookNote;
|
||
try {
|
||
const cfg = JSON.parse(readFileSync(file, 'utf8'));
|
||
const entries = cfg?.hooks?.PermissionRequest;
|
||
hookRegistered = Array.isArray(entries) && entries.some(e => Array.isArray(e?.hooks) && e.hooks.length > 0);
|
||
hookNote = hookRegistered ? `钩子已注册(${file})` : `钩子清单里没有 PermissionRequest(${file})`;
|
||
} catch (e) {
|
||
hookNote = `读不到钩子清单(${file}):${e?.message || e}`;
|
||
}
|
||
|
||
if (gateReady) {
|
||
return {
|
||
enforcement: 'native',
|
||
reason:
|
||
`执行类动作只能经我们自己的门禁(run_command / write_file 逐次请示),` +
|
||
`平台自带危险工具已禁用 ${denied.length} 项(--mode ${mode})` +
|
||
(hookRegistered ? ';桌面模式另有平台钩子' : '')
|
||
};
|
||
}
|
||
if (hookRegistered && modeReachesPermissionHook(mode)) {
|
||
return { enforcement: 'native', reason: hookNote };
|
||
}
|
||
return {
|
||
enforcement: 'advisory',
|
||
reason: hookRegistered ? `钩子已注册但不参与(--mode ${mode})` : `${hookNote},且门禁自检未通过`
|
||
};
|
||
}
|
||
|
||
/** 描述错误:把「网关可达但返回 4xx」与「连不上」分开 —— 两者的应对完全不同。 */
|
||
export function describeError(e) {
|
||
if (e instanceof GatewayError) {
|
||
return `网关返回 HTTP ${e.status}(${e.path}):${typeof e.body === 'string' ? e.body : e.message}`;
|
||
}
|
||
return e?.message || String(e);
|
||
}
|
||
|
||
/**
|
||
* 把 ZCode 的事件流里「值得进日志」的那几条提出来。
|
||
*
|
||
* 一轮实测能吐 6893 条事件,全记等于没有日志。只记两类:
|
||
* **工具调用**(模型在干什么、有没有触发授权)与**错误**。
|
||
* 邮件驱动的会话没有界面,这两类是唯一能回答
|
||
* 「它是不是卡住了 / 为什么一直没有授权询问」的信息。
|
||
*/
|
||
export function describeRunEvent(event) {
|
||
const t = String(event?.type || '');
|
||
if (t === 'tool.call.started') return `工具 ${event.toolName || '?'}`;
|
||
if (t.startsWith('tool.permission')) {
|
||
return `权限 ${event.decision || event.behavior || '?'}(${event.toolName || event.ruleId || '?'})`;
|
||
}
|
||
if (/error|failed|denied|aborted/i.test(t)) return `事件 ${t}`;
|
||
return null;
|
||
}
|
||
|
||
/**
|
||
* 造一个驱动实例。
|
||
*
|
||
* 依赖全部注入,所以整条流水线(事件 → 提示词 → 一轮 → 回信判定 → 发信载荷)
|
||
* 可以在没有模型、没有 ZCode 的情况下被端到端断言。
|
||
*/
|
||
export function createDriver({ client, runTurnFn = runTurn, logFn = log, env = process.env, config } = {}) {
|
||
// 配置可覆盖:测试需要把工作目录指到临时目录,不能碰真实的家目录。
|
||
const CFG = { ...CONFIG, ...(config || {}) };
|
||
const delivered = new BoundedSet(MAX_TRACKED_MAILS);
|
||
/** AgentMail 会话 id → { zcodeSessionId, cwd, tier, turns } */
|
||
const sessions = new BoundedMap(MAX_TRACKED_SESSIONS);
|
||
const queue = [];
|
||
let running = false;
|
||
/** 当前在途回合的杀进程函数(关停时要终止它,否则会留下跑工具的孤儿)。 */
|
||
let currentKill = null;
|
||
|
||
/** 注册表只用于日志与自检:它让「为什么一轮授权询问都没发生」有据可查。 */
|
||
const stats = { turns: 0, relays: 0, skippedRelay: 0, failures: 0 };
|
||
|
||
function resolveCwd(data) {
|
||
const sessionId = data?.session_id || data?.mail_id || 'unknown';
|
||
const fallback = zcodeSessionFallback(sessionId, CFG.workspaceRoot);
|
||
const { cwd, grouped } = resolveWorkspaceCwd(data?.to_workspace, fallback);
|
||
if (!grouped) logFn(`会话 ${sessionId} 没有可用的 to_workspace,用兜底目录 ${cwd}`);
|
||
ensureCwd(cwd, grouped);
|
||
return cwd;
|
||
}
|
||
|
||
async function relay({ data, text, kind }) {
|
||
const fromHuman = data?.from_human === true;
|
||
const decision = autoRelayDecision({ fromHuman, replyTo: data?.from_name });
|
||
if (!decision.relay) {
|
||
logFn(`不自动转发(${decision.reason})`);
|
||
stats.skippedRelay++;
|
||
return false;
|
||
}
|
||
|
||
const sessionId = data?.session_id || '';
|
||
const sent = readExplicitSends(explicitSendsFile(env), {
|
||
sessionId,
|
||
// 只认本轮之后的记录:早于本轮的发信属于上一次往返,不该让这一轮沉默。
|
||
since: Date.now() - CFG.turnTimeoutMs
|
||
});
|
||
if (shouldSkipAutoRelay(sent, data.from_name, data.mail_id)) {
|
||
logFn(`本轮模型已主动回信 ${data.from_name},跳过自动转发`);
|
||
stats.skippedRelay++;
|
||
return false;
|
||
}
|
||
|
||
// relay + relay_key 走免配额通道:模型已经把话说完了,驱动只是把它搬进邮件。
|
||
// 对搬运收配额会让「配额用尽」变成「连交代都做不到」。
|
||
const relayKey = clampRelayKey(`zcode:${data.mail_id || kind}`);
|
||
await client.post('/mail/send', {
|
||
to: data.from_name,
|
||
subject: replySubject(data.subject),
|
||
body: text,
|
||
reply_to: data.mail_id || '',
|
||
relay: 'summary',
|
||
relay_key: relayKey
|
||
});
|
||
stats.relays++;
|
||
logFn(`已回信给 ${data.from_name}(${text.length} 字)`);
|
||
return true;
|
||
}
|
||
|
||
async function processMail(data) {
|
||
const sessionId = data?.session_id || '';
|
||
const tier = normalizeMode(data?.permission_mode);
|
||
const mode = zcodeModeForTier(tier);
|
||
// 危险的自带工具一律拿掉(三个档位同一张审过的清单);需要动手时,
|
||
// 模型改用我们自己的 run_command / write_file,门禁就在那里面。
|
||
const disallowedTools = denylistForTier(tier);
|
||
const cwd = resolveCwd(data);
|
||
const prev = sessions.get(sessionId);
|
||
const resume = prev?.zcodeSessionId || '';
|
||
|
||
logFn(`处理 ${data.mail_id}|${describeTier(tier, mode)}|cwd=${cwd}${resume ? `|续会话 ${resume}` : ''}`);
|
||
logFn(`已禁用 ${disallowedTools.length} 个自带工具(Bash/Write/Edit/js/…),执行类动作走 AgentMail 门禁`);
|
||
if (modeReachesPermissionHook(mode)) {
|
||
// 平台会问、我们的钩子会转达 —— 说明档位映射被配置改回了 build/edit。
|
||
logFn(`注意:--mode ${mode} 下平台自带工具会产生权限询问(映射被覆盖过?)`);
|
||
}
|
||
|
||
const prompt = buildMailPrompt({ agentName: CFG.agentName, data });
|
||
|
||
const outcome = await runTurnFn(
|
||
{
|
||
prompt,
|
||
cwd,
|
||
mode,
|
||
maxTurns: CFG.maxTurns,
|
||
resumeSessionId: resume || undefined,
|
||
disallowedTools,
|
||
turnTimeoutMs: CFG.turnTimeoutMs,
|
||
cliPath: CFG.cliPath,
|
||
// 注入给 ZCode 进程(→ 继承给插件、钩子、MCP 服务器):
|
||
// 授权钩子靠 AGENTMAIL_SESSION_ID 判断「有没有本地界面」,
|
||
// 靠 AGENTMAIL_PERMISSION_MODE 决定档位。
|
||
env: {
|
||
AGENTMAIL_SESSION_ID: sessionId,
|
||
AGENTMAIL_PERMISSION_MODE: tier,
|
||
AGENTMAIL_MAIL_SUBJECT: data?.subject || '',
|
||
AGENTMAIL_REPLY_TO: data?.mail_id || ''
|
||
},
|
||
onEvent: event => {
|
||
const line = describeRunEvent(event);
|
||
if (line) logFn(line);
|
||
}
|
||
},
|
||
{
|
||
log: logFn,
|
||
onChild: kill => {
|
||
currentKill = kill;
|
||
}
|
||
}
|
||
);
|
||
currentKill = null;
|
||
|
||
stats.turns++;
|
||
|
||
if (outcome.sessionId) {
|
||
sessions.set(sessionId, {
|
||
zcodeSessionId: outcome.sessionId,
|
||
cwd,
|
||
tier,
|
||
turns: (prev?.turns || 0) + 1
|
||
});
|
||
}
|
||
|
||
const failed = outcome.timedOut || (outcome.exitCode !== 0 && !outcome.response);
|
||
if (failed) {
|
||
stats.failures++;
|
||
const reason = outcome.timedOut
|
||
? `回合超时(${Math.round(CFG.turnTimeoutMs / 1000)} 秒),已终止进程树`
|
||
: `ZCode 退出码 ${outcome.exitCode}${outcome.stderrTail ? `:\n${outcome.stderrTail}` : ''}`;
|
||
logFn(`一轮失败:${reason}`);
|
||
// 失败必须回信:否则发件人只看到「信发出去了,然后再无音讯」。
|
||
try {
|
||
await client.post('/mail/send', {
|
||
to: data?.from_name,
|
||
subject: `处理失败: ${data?.subject || '(无主题)'}`,
|
||
body: renderTurnFailure([{ kind: outcome.timedOut ? '超时' : 'CLI 失败', error: reason }], data?.subject),
|
||
reply_to: data?.mail_id || '',
|
||
relay: 'summary',
|
||
relay_key: clampRelayKey(`zcode-failure:${data?.mail_id || sessionId}`)
|
||
});
|
||
} catch (e) {
|
||
logFn(`失败回报也发不出去:${describeError(e)}`);
|
||
}
|
||
return { ok: false, reason };
|
||
}
|
||
|
||
const text = String(outcome.response || '').trim();
|
||
if (!text) {
|
||
// 退出码 0 但没有最终文本:常见于模型只调了工具就结束。
|
||
// 这时**不冒充**回信(会让收件人以为模型什么都没做),但要留下日志。
|
||
logFn('这一轮没有产出最终文本,不自动回信(若模型自己发过信,那封就是答复)');
|
||
return { ok: true, relayed: false };
|
||
}
|
||
|
||
return { ok: true, relayed: await relay({ data, text, kind: 'mail' }) };
|
||
}
|
||
|
||
async function drain() {
|
||
if (running) return;
|
||
running = true;
|
||
try {
|
||
while (queue.length) {
|
||
const data = queue.shift();
|
||
try {
|
||
await processMail(data);
|
||
} catch (e) {
|
||
// 一封邮件处理崩了不能把驱动带走:后面还有很多信。
|
||
logFn(`处理 ${data?.mail_id} 时异常:${describeError(e)}`);
|
||
}
|
||
}
|
||
} finally {
|
||
running = false;
|
||
}
|
||
}
|
||
|
||
/** SSE 事件入口。必须廉价 —— 它跑在读循环上。 */
|
||
function handleEvent(type, data) {
|
||
if (type !== 'new_mail') return;
|
||
if (data?.role && data.role !== 'to' && data.role !== 'cc') return;
|
||
const id = data?.mail_id;
|
||
if (!id || delivered.has(id)) return;
|
||
delivered.add(id);
|
||
queue.push(data);
|
||
void drain();
|
||
}
|
||
|
||
async function catchUp(pendingMails) {
|
||
const mails = selectCatchup(pendingMails, delivered);
|
||
if (!mails.length) return 0;
|
||
logFn(`补投 ${mails.length} 封停机期间到达的邮件`);
|
||
for (const m of mails) handleEvent('new_mail', m.data ?? m);
|
||
return mails.length;
|
||
}
|
||
|
||
return {
|
||
handleEvent,
|
||
catchUp,
|
||
processMail,
|
||
stats,
|
||
sessions,
|
||
delivered,
|
||
/**
|
||
* 关停:终止在途回合。
|
||
*
|
||
* 不做这件事的后果是——systemd 杀掉驱动之后,那个 ZCode 进程还在跑工具,
|
||
* 而既没有驱动看着它,也没有本地界面看着它。宁可丢掉这一轮的工作。
|
||
*/
|
||
abort() {
|
||
// 先取后清:杀过就算完,重复关停(SIGTERM 后再来一个)不该重复杀。
|
||
const kill = currentKill;
|
||
currentKill = null;
|
||
if (kill) {
|
||
logFn('关停:终止在途的 ZCode 回合');
|
||
try {
|
||
kill('SIGTERM');
|
||
} catch {
|
||
/* 已经结束了 */
|
||
}
|
||
}
|
||
queue.length = 0;
|
||
}
|
||
};
|
||
}
|
||
|
||
// ─── 真实入口 ───────────────────────────────────────────────────────
|
||
|
||
async function main() {
|
||
const client = new GatewayClient(process.env);
|
||
const missing = client.checkConfig();
|
||
if (missing.length) {
|
||
log(`配置不完整,缺少 ${missing.join('、')};驱动不会启动(静默启动会让信永远没人处理)`);
|
||
process.exit(1);
|
||
}
|
||
|
||
const driver = createDriver({ client, logFn: log });
|
||
let caughtUp = false;
|
||
let timer;
|
||
|
||
try {
|
||
await client.register();
|
||
log(`已接入 ${client.baseURL},身份 ${client.agentName}`);
|
||
} catch (e) {
|
||
// 密钥未登记时说清该做什么,别只留一句 401。
|
||
log(`注册失败:${describeError(e)}`);
|
||
log('若提示密钥无效,请让管理员在 AgentMail 后台登记这把密钥。');
|
||
}
|
||
|
||
const beat = async () => {
|
||
try {
|
||
// mode_enforcement 只声明得出来的事:挡得住工具的是插件里的授权钩子,
|
||
// 钩子没注册时我们什么也拦不住(见 detectModeEnforcement)。
|
||
const res = await client.post('/agent/heartbeat', {
|
||
mode_enforcement: detectModeEnforcement().enforcement
|
||
});
|
||
if (!caughtUp) {
|
||
caughtUp = true;
|
||
await driver.catchUp(res?.pending_mails);
|
||
}
|
||
} catch {
|
||
// 心跳失败不刷错误日志:真连不上时网关会把它判成离线,那才是可见信号。
|
||
}
|
||
};
|
||
await beat();
|
||
timer = setInterval(beat, 30_000);
|
||
|
||
createSSEClient({
|
||
authHeaders: () => client.authHeaders(),
|
||
baseURL: client.baseURL,
|
||
path: '/api/v1/events/stream',
|
||
log,
|
||
onEvent: (type, data) => driver.handleEvent(type, data)
|
||
});
|
||
// 留住句柄:关停时要真的断开,否则重连定时器还在跑(进程虽然马上就退,
|
||
// 但那是侥幸而不是设计)。
|
||
const sse = createSSEClient({
|
||
authHeaders: () => client.authHeaders(),
|
||
baseURL: client.baseURL,
|
||
path: '/api/v1/events/stream',
|
||
log,
|
||
onEvent: (type, data) => driver.handleEvent(type, data)
|
||
});
|
||
|
||
const shutdown = reason => {
|
||
log(`收到 ${reason},关停中…(已处理 ${driver.stats.turns} 轮,回信 ${driver.stats.relays} 封)`);
|
||
if (timer) clearInterval(timer);
|
||
driver.abort();
|
||
sse.stop();
|
||
// 给杀进程留一点时间再退:自己先死会把 ZCode 变成孤儿。
|
||
setTimeout(() => process.exit(0), 1200);
|
||
};
|
||
for (const sig of ['SIGINT', 'SIGTERM']) process.on(sig, () => shutdown(sig));
|
||
|
||
const enforcement = detectModeEnforcement();
|
||
log(`档位强制力自报:${enforcement.enforcement}(${enforcement.reason})`);
|
||
log(`驱动就绪:CLI ${basename(CONFIG.cliPath)},回合上限 ${Math.round(CONFIG.turnTimeoutMs / 1000)} 秒`);
|
||
}
|
||
|
||
// 直接执行时启动;被 import 时只导出(测试要用 createDriver)。
|
||
//
|
||
// 判断**必须解析软链**(见 lib/is-main.mjs):生产布局是 `current` 软链,
|
||
// 直接比字符串会让驱动什么都不做就退出(退出码 0)—— 而那看起来
|
||
// 与「服务正常启动但没收到信」一模一样。
|
||
if (isMainModule(import.meta.url)) {
|
||
main().catch(e => {
|
||
log(`启动失败:${describeError(e)}`);
|
||
process.exit(1);
|
||
});
|
||
}
|