Files
MailUI4Agents/plugins/zcode-mail-bridge/lib/sse-client.js
JianFeeeee c774904c0c feat(zcode): 授权桥 —— PermissionRequest 钩子把危险工具授权交给人
第二步:让 ZCode 上的 Bash/Write/Edit 授权走 AgentMail 的人工审批,
而不是只靠本地界面。

钩子契约从 CLI 产物里逆出来(不猜协议):
- 输入走 stdin:{hook_event_name, tool_name, tool_input, session_id, permission_mode…}
- 输出走 stdout,schema **严格**:{"decision":"approve"} / {"decision":"block","reason"}
  多一个键就会报 "Hook stdout failed HookJSONOutput schema validation"
- 空输出 / 不以 { 开头 = 不表态;exit 2 = 拒绝;其它非零 = 钩子失败
- 注入的环境变量含 ZCODE_PLUGIN_ROOT / ZCODE_PLUGIN_DATA / ZCODE_SESSION_ID
  (MCP 配置里用 ZCODE_SESSION_ID 反而会抛「需要运行时会话上下文」)

档位判定与 pi 桥逐条对齐(plan 直接拒 / workspace 问人 / full 批准),
判定逻辑抽成纯函数 lib/hook-policy.mjs 以便穷举:
其中 full 档必须**返回批准而不是不表态** —— 钩子一旦触发说明 ZCode 本会去问人,
不表态等于让那个询问照常发生,full 档就退化成了 workspace 档。

钩子自己开 SSE 等决定,不依赖桥进程:网关的 SSE 是扇出的
(clients 按唯一 id 存,SendToAgent 推给该 Agent 的所有客户端),
一次性进程也能订阅到自己那条 permission_decision。这样交互模式下同样可用
(人自己开着 ZCode 干活时并没有桥在跑)。先建连再发请求是有意的:
反过来会有一个窗口,人在窗口内点的同意推送给当时还不存在的客户端。

fail closed 但区分模式:永久失败(409/4xx)一律拒绝;暂时失败在
AGENTMAIL_SESSION_ID 非空(邮件驱动、没有本地界面兜底)时拒绝,
交互模式则不表态让人就地决定。

「一直同意」落盘(lib/grants-file.mjs):钩子是一个事件一个进程,
不落盘那个选项就是骗人的。判定仍交给共用的 permission-grants.js。

共用模块同源范围扩到 9 个(新增 permission-mode / relay-key /
permission-grants / sse-client)—— 档位语义与决策判定分叉会让「同意」
在 ZCode 上悄悄变成另一种意思。

验证:
- 单元 229/229(新增 hook-policy 14 项、grants-file 8 项,含反向对照)
- 共用模块四方同源检查通过
- 授权桥端到端 5/5,全部带反向对照:
  同意→approve;拒绝→block 且原因必须来自人的拒绝(不能是超时兜底);
  plan 档拒绝且**不产生**任何权限邮件;无人可问(409)→fail closed;
  非守卫工具→不表态
- `zcode plugins list` → agentmail@inline [enabled],hooks: 1,
  mcp: plugin:agentmail:agentmail

我自己写错的两处判据(都已修,值得记下):
1. 待决权限列表里有历史积压(实测 6 条,含其它 Agent 的条目),
   只按「第一条新的」取会拿到无关请求 —— 于是人点了同意而钩子在等自己那条,
   最后超时。第一版还把这个超时误报成「拒绝路径通过」。
   现在按「启动前快照差集 + session_id + agent_name」三重过滤。
2. 「无人可问」控制组最初传了个非 UUID 的 session id,走的是 400(参数错),
   验不到 409 那条真实路径。改为真的造一条只有 Agent 没有人类的会话。
2026-09-12 14:09:10 +08:00

180 lines
6.1 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.

/**
* 共用 SSE 客户端:跨 TCP 分片保帧状态 + Last-Event-ID 断点续传。
*
* 三平台桥原本各自手写 SSE 解析,且都有同一个 bug
* - `evt` / `data` 是每次 `read()` 的局部变量TCP 把一帧
* `event: xxx\ndata: {...}\n\n` 切在换行处时,第一段只剩
* `event:` 而第二段只有 `data:` —— 整帧被静默丢弃。
* - 重连不带 `Last-Event-ID`,断线期间的事件只在服务端环形
* 缓冲里等着永远回放不出来Gateway 有 per-agent ring buffer
* pi 与 homeagent 已正确利用DSH/opencode 没有)。
*
* 这个模块把 pi 桥 `src/gateway.mjs` 里那份验证过的实现抽成共用件,
* 三桥逐字节同源deploy/check-shared-libs.sh 校验)。
*/
/**
* 增量 SSE 帧解析器。
*
* `push(chunk)` 可以喂任意切分的文本片段,返回本次完整解析出的事件数组。
* 所有跨帧状态(缓冲、当前 event/data/id都保存在闭包里**不随 chunk 重置** ——
* 这正是原实现丢帧的根因。
*
* 协议细节:
* - `:` 开头 = 注释/心跳,忽略
* - `id:` / `event:` / `data:` 各取字段;`data:` 后的单个空格是分隔符
* - 多行 data 用 `\n` 拼接
* - 空行 = 帧结束;只有 event 与 data 都非空才派发(与旧行为一致)
* - 兼容 CRLF
* - `lastEventId` 在**派发之前**记下:回调抛异常也不该让断点回退。
*
* @returns {{push: (chunk: string) => Array<{event: string, data: string, id: string}>,
* reset: () => void,
* lastEventId: () => string,
* setLastEventId: (id: string) => void}}
*/
export function createFrameParser() {
let buffer = '';
let lastEventId = '';
let curEvent = '';
let curData = '';
let curId = '';
function push(chunk) {
buffer += chunk;
const events = [];
const lines = buffer.split('\n');
// 最后一段可能是被切断的半行,留到下一个 chunk
buffer = lines.pop() ?? '';
for (let line of lines) {
if (line.length > 0 && line.charAt(line.length - 1) === '\r') {
line = line.slice(0, -1);
}
if (line.startsWith(':')) continue;
if (line.startsWith('id:')) {
curId = line.slice(3).trim();
} else if (line.startsWith('event:')) {
curEvent = line.slice(6).trim();
} else if (line.startsWith('data:')) {
let value = line.slice(5);
if (value.startsWith(' ')) value = value.slice(1);
curData = curData.length > 0 ? `${curData}\n${value}` : value;
} else if (line === '') {
if (curEvent.length > 0 && curData.length > 0) {
if (curId.length > 0) lastEventId = curId;
events.push({ event: curEvent, data: curData, id: curId });
}
curEvent = '';
curData = '';
curId = '';
}
}
return events;
}
function reset() {
buffer = '';
curEvent = '';
curData = '';
curId = '';
}
return {
push,
reset,
lastEventId: () => lastEventId,
setLastEventId: (id) => { lastEventId = id || ''; },
};
}
/**
* @param {object} deps
* @param {() => Record<string,string>} deps.authHeaders 认证头(每次重连重取,密钥可能已换)
* @param {string} deps.baseURL Gateway 基地址(不带末尾 /
* @param {string} deps.path SSE 路径(如 /api/v1/events/stream
* @param {(evt: string, data: any) => void} deps.onEvent 事件分发回调
* @param {(msg: string) => void} [deps.log] 日志回调(默认 console.error
* @returns {{stop: () => void}} stop() 终止重连与在途请求
*/
export function createSSEClient({ authHeaders, baseURL, path, onEvent, log = console.error }) {
const controller = new AbortController();
const parser = createFrameParser();
function stop() {
controller.abort();
}
function reconnect(delay) {
if (controller.signal.aborted) return;
setTimeout(() => connect(), delay);
}
function connect() {
if (controller.signal.aborted) return;
const headers = { ...authHeaders(), Accept: 'text/event-stream' };
// 只有 lastEventId 非空(= 已经收过事件)时才是重连:首次连接不带,
// 否则服务端会把环形缓冲里的旧事件全回放一遍,插件重启后重复处理一批已处理的邮件。
const lastEventID = parser.lastEventId();
if (lastEventID) {
headers['Last-Event-ID'] = lastEventID;
log(`SSE 重连,从事件 ${lastEventID} 之后续传`);
}
fetch(`${baseURL}${path}`, { headers, signal: controller.signal })
.then((res) => {
if (!res.ok || !res.body) {
log(`SSE 建连失败: HTTP ${res.status}`);
return reconnect(5000);
}
const reader = res.body.getReader();
const decoder = new TextDecoder();
function read() {
reader.read().then(({ done, value }) => {
if (done) {
parser.reset();
return reconnect(3000);
}
for (const ev of parser.push(decoder.decode(value, { stream: true }))) {
try {
onEvent(ev.event, JSON.parse(ev.data));
} catch (e) {
log(`SSE 事件处理失败: ${e?.message || e}`);
}
}
read();
}).catch((e) => {
if (controller.signal.aborted) return;
log(`SSE 读取中断: ${e?.message || e}`);
parser.reset();
reconnect(5000);
});
}
read();
})
.catch((e) => {
if (controller.signal.aborted) return;
log(`SSE 连接错误: ${e?.message || e}`);
parser.reset();
reconnect(5000);
});
}
connect();
return {
stop,
/**
* 清掉断点(不终止连接)。
*
* 换 Gateway 地址时必须调lastEventID 是**旧** Gateway 环形缓冲里的序号,
* 拿去问新 Gateway 会命中一段完全无关的历史(或直接被拒),
* 得到的事件属于别人的会话。
*/
reset: () => parser.setLastEventId(''),
};
}