# 发生了什么 pi-lens 内置「安全格式化」:它会自动安装 biome 并对**编辑过的文件**跑 `biome format --write`。本机原先没有任何 biome 配置,于是 biome 用它自己的 默认值 —— tab 缩进 + 双引号 —— 把文件整体重写。 我在19a3161那次提交里用了 `git add -A`,把这批与功能无关的重排一起扫了进去: 约 7000 行改动散落在 20 个文件上,使那次提交无法审查,还掩盖了 server/internal/handler/permission.go 的一处删行(实为文件末尾空行,无代码丢失)。 # 为什么是「关掉」而不是「配置成我们的风格」 试过把缩进/引号/lineWidth 全部对齐本仓库习惯(biome.json + space/2/single/ lineWidth 120):`biome format --write` 仍然改动 17 个文件。原因是本仓库从未按 biome 的规则排版过 —— 注释按语义换行、数组与调用按可读性手工折行, 这些无法由格式化器还原。也就是说只要格式化器开着,每次编辑都会产生与内容无关的 大面积 diff,把真正的改动埋掉。 因此 biome.jsonc 里 formatter 与 linter 都关闭:本仓库的静态检查由 tsc / go vet / tree-sitter / ast-grep 与各自测试套件承担,不引入会改动无关行的 自动修复。 (pi-lens 这一版把 format 服务的 enabled 硬编码为 true,没有配置开关, 所以只能在仓库侧用 biome 配置让它不动文件;已验证 `biome format --write` 对这些文件零改动。) # 本提交内容 把19a3161里除「有意改动」外的 20 个文件还原到重排前的样子。19a3161中真正有意的改动是 deploy/install.sh 的扩展注册与 plugins/pi-mail-bridge/extension/index.ts 新文件,两者原样保留。 验证:Go 全量、三桥插件(320/362/409)、前端 196 全绿; `biome format --write` 对还原后的文件零改动。
180 lines
6.1 KiB
JavaScript
180 lines
6.1 KiB
JavaScript
/**
|
||
* 共用 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(''),
|
||
};
|
||
}
|