Files
MailUI4Agents/plugins/pi-mail-bridge/src/session-scan.mjs
JianFeeeee 79c4171c9d feat: L0 线协议冻结 + 附件链路修复 + 人/Agent 区分
L0 核心:
- 严格解码 Decode(DisallowUnknownFields) 全覆盖 29 个 DecodeBody 调用点
- DecodeLenient 心跳专用:容忍新字段但回报 unknown_fields
- 400 消息列出本端点接受的全部字段(jsonFieldNames 反射 tag)
- 日历 status 校验(create 补字段 + update 拦非法值)
- 新增 strictdecode_test.go 10 例 + blob/list_test.go 6 例

A-4 附件挂载回滚:checkAttachable 在 CreateMail 前校验,失败按
解挂→释放 relay→删邮件→退预算回滚,幽灵邮件这条路堵住了

A-5 反向 GC:blob.Store.List() 枚举磁盘(跳 .upload-*),
SweepUnreferencedBlobs 按 attachments + calendar_attachments 反查,
48h 年龄下限兜上传窗口。已接进每小时 sweep 循环

C 人/Agent 区分:四个读路径 + threadCols 补 from_human / to_human
(EXISTS users 判定),models.Mail 加 ToHuman。前端判据从
workspace 启发式改成显式布尔,mailCounterpart/sessionCounterpart
从 session_workspace 取 path(修 dsh@dsh 拼接 bug)

契约文档:SSE new_mail 补 4 字段(in_reply_to/from_human/
permission_mode/permission_enforcement),B-5 加 B-5.6
(Agent→Agent 不转发),B-3.4 MUST 改条件式,心跳补 mode_enforcement
+ unknown_fields,demo 死链修复 + from_human 检查
验收清单加 Agent→Agent 负向对照项
2026-09-06 15:18:06 +08:00

322 lines
12 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 会话目录的增量扫描器 —— 替换心跳路径上的 `SessionManager.listAll()`。
*
* # 为什么不能用 listAll
*
* 心跳每 30 秒要上报一次平台会话快照(`snapshotPiSessions`),而它只用到四个
* 字段:`id` / `cwd` / `name` / `modified`。`SessionManager.listAll()` 为了拿到
* 这四个字段,把 `~/.pi/agent/sessions` 下**每个 .jsonl 的每一行**都读进来并
* `JSON.parse`,还顺手把所有消息正文拼成一个 `allMessagesText` 大字符串。
*
* 本机实测(115 个文件 / 145MB,其中单个会话文件 29MB、单行最长 2.6MB):
*
* listAll() 1431ms RSS 41 → 323MB(heapUsed 141MB)
* 仅读 header 3ms RSS 41 → 46MB
* 本模块(冷启动) 516ms RSS 41 → 131MB
* 本模块(稳态) 3ms 重扫 0 字节
*
* 那 282MB 每 30 秒分配一次、随即变成垃圾。GC 收得掉(所以 RSS 呈锯齿而不是
* 单调上升),但代价是:常驻内存被垃圾撑到 300MB 上下,且每拍有 1.4 秒的同步
* 解析跑在事件循环上 —— 那期间 SSE 读循环停着,新邮件事件在 TCP 缓冲区排队。
*
* # 三条省法
*
* 1. **id / cwd / created 只在首行**。header 是第一行,读 4KB 就够,不必读全文。
* 2. **name 来自 `session_info` 行**,而那种行只有几百字节。按行扫描时长度超过
* 上限的行**直接跳过、不materialize**,于是 2.6MB 的 message 行不进内存。
* 3. **文件是 append-only 的**。缓存 `size`,下一拍只扫 `[上次 size, 现 size)`
* 这段尾巴 —— 没有新消息的会话一个字节都不读。
*
* `modified` 改用 `stat.mtime`:listAll 是从最后一条消息的活动时间算的,
* 而快照只拿它排序(「最近在谈的排前面」),mtime 表达的正是这件事,且免费。
*
* # 为什么这个模块不进 lib/(不与另两个平台共用)
*
* 它读的是 pi 自己的磁盘格式。opencode 的 `session.list()` 是进程内调用,
* DSH 走 `sessionQuery` 的语料观测 —— 两者都没有「扫目录读文件」这一步,
* 强行抽象成共用模块只会得到一个谁都不合身的接口。
*/
import { createReadStream } from 'node:fs';
import { readdir, stat, open } from 'node:fs/promises';
import { join } from 'node:path';
import { StringDecoder } from 'node:string_decoder';
/**
* 单行长度上限(字符)。超过这个长度的行不参与解析。
*
* `session_info` 行的构成是固定的:type + 两个 8 字节 id + ISO 时间戳 + name,
* 而 pi-web 的标题生成器把 name 截到 60 字符。几百字节封顶,留 4096 是十倍余量。
*
* 这个上限同时是**内存上界**:扫描时跨块累积的未完成行一旦超过它就被丢弃,
* 因此无论会话里有多大的一条 message(实测见过 2.6MB),扫描峰值都不受影响。
*/
const MAX_LINE_CHARS = 4096;
/** header 只读这么多字节。第一行是 `{"type":"session",...}`,远不到 4KB。 */
const HEADER_READ_BYTES = 4096;
/** `session_info` 行的判别串。先做子串命中再 JSON.parse,省掉绝大多数解析。 */
const SESSION_INFO_NEEDLE = '"type":"session_info"';
/**
* 读会话文件的 header(第一行)。
*
* @param {string} file
* @returns {Promise<{id: string, cwd: string, created: Date} | null>}
* 不是合法会话文件时返回 null(首行不是 session 类型、空文件、读不动)。
*/
async function readHeader(file) {
let fh;
try {
fh = await open(file, 'r');
const buf = Buffer.allocUnsafe(HEADER_READ_BYTES);
const { bytesRead } = await fh.read(buf, 0, HEADER_READ_BYTES, 0);
if (bytesRead === 0) return null;
const text = buf.subarray(0, bytesRead).toString('utf8');
const nl = text.indexOf('\n');
// 没有换行说明首行比 4KB 还长 —— 那不是 header(header 是固定几个字段)。
if (nl < 0) return null;
const entry = JSON.parse(text.slice(0, nl));
if (entry?.type !== 'session' || typeof entry.id !== 'string') return null;
const ts = typeof entry.timestamp === 'string' ? new Date(entry.timestamp) : null;
return {
id: entry.id,
// 老会话的 cwd 是空串(pi 的 SessionInfo 注释写明了),照实传下去 ——
// snapshotPiSessions 会按空 workspace 上报,不该拿桥自己的 cwd 冒充。
cwd: typeof entry.cwd === 'string' ? entry.cwd : '',
created: ts && !Number.isNaN(ts.getTime()) ? ts : new Date(0),
};
} catch {
// 读不动、JSON 坏了、文件正好被删 —— 一律当「不是会话」。
// 单个坏文件不该让整份快照失败(listAll 也是这个策略)。
return null;
} finally {
await fh?.close().catch(() => {});
}
}
/**
* 扫一段字节区间,返回其中**最后一个** `session_info` 的 name。
*
* 「最后一个」而不是第一个:会话可以被改名多次,也可以显式清名
* (`session_info` 不带 name = 清掉)。语义与 SDK 的 buildSessionInfo 一致 ——
* 最新的那条生效。
*
* @param {string} file
* @param {number} from 起始字节(含)。append-only 文件里它一定落在行首。
* @param {number} to 结束字节(不含)
* @returns {Promise<{name: string | undefined, found: boolean}>}
* `found` 为假表示这段里没有任何 session_info —— 调用方应保留上一次的 name,
* 而不是把它当成「被清空了」。
*/
async function scanRangeForName(file, from, to) {
if (to <= from) return { name: undefined, found: false };
let name;
let found = false;
const decoder = new StringDecoder('utf8');
let pending = '';
// 当前这一行已经超过上限 → 丢弃它剩下的部分,直到下一个换行。
// 这是内存上界的实现:巨大的 message 行永远不会被拼出来。
let skipping = false;
const consider = (line) => {
if (line.length > MAX_LINE_CHARS) return;
if (!line.includes(SESSION_INFO_NEEDLE)) return;
try {
const entry = JSON.parse(line);
if (entry?.type !== 'session_info') return;
found = true;
const n = typeof entry.name === 'string' ? entry.name.trim() : '';
name = n || undefined;
} catch {
// 半截行(起点没对齐、文件正在被写)解析失败 —— 忽略即可,
// 下一拍尾巴长出来之后会重新看到完整的那一行。
}
};
const stream = createReadStream(file, { start: from, end: to - 1 });
for await (const chunk of stream) {
const text = decoder.write(chunk);
let start = 0;
for (;;) {
const nl = text.indexOf('\n', start);
if (nl < 0) break;
if (!skipping) consider(pending + text.slice(start, nl));
pending = '';
skipping = false;
start = nl + 1;
}
if (skipping) continue;
pending += text.slice(start);
if (pending.length > MAX_LINE_CHARS) {
pending = '';
skipping = true;
}
}
const tail = decoder.end();
if (!skipping) {
pending += tail;
// 末行没有换行符时也要看一眼(正在被写入的会话就是这种状态)
if (pending) consider(pending);
}
return { name, found };
}
/**
* 建一个扫描器。缓存跨拍存活,因此要在插件启动时建一次、之后复用。
*
* @param {object} [opts]
* @param {string} [opts.sessionsDir] 会话根目录。默认 `~/.pi/agent/sessions`
* (由调用方传 `join(getAgentDir(), 'sessions')`,这里不 import SDK
* —— 让这个模块可以脱离 SDK 单测)。
* @returns {{scan: () => Promise<object[]>, stats: () => object}}
*/
export function createSessionScanner({ sessionsDir } = {}) {
if (!sessionsDir) throw new Error('createSessionScanner 需要 sessionsDir');
/**
* file -> { size, id, cwd, created, name, hasName }
*
* 这张表的键是磁盘上真实存在的文件,每次 scan 都会把消失的文件删掉 ——
* 所以它不需要额外的上界:会话文件被删(pi 侧清理历史)时条目跟着走。
*
* `hasName` 与 `name === undefined` 不同:前者是「曾经见过 session_info」,
* 后者可能是「见过但被清名了」。区分它们才能让增量扫描保留上一次的结论。
*/
const cache = new Map();
let fullScans = 0;
let tailScans = 0;
let tailBytes = 0;
async function collectFiles() {
const out = [];
let dirs;
try {
dirs = await readdir(sessionsDir, { withFileTypes: true });
} catch (e) {
// 目录不存在(pi 从没跑过)= 确实一条会话都没有,空列表是正确答案。
//
// 其他错误(权限、I/O)必须**抛出去**:调用方据此省略 platform_sessions
// 字段,保留服务端镜像。返回空数组的语义是「平台确实没有会话」,
// 会把镜像抹掉(W-3 / N-7)—— 一次 EACCES 就能清空别人的补全候选。
if (e?.code === 'ENOENT') return out;
throw e;
}
for (const d of dirs) {
if (!d.isDirectory() && !d.isSymbolicLink()) continue;
const dir = join(sessionsDir, d.name);
try {
for (const f of await readdir(dir)) {
if (f.endsWith('.jsonl')) out.push(join(dir, f));
}
} catch {
// 单个子目录读不动(权限、正被删)不该拖垮整轮
}
}
return out;
}
/**
* 扫一遍,返回与 `SessionManager.listAll()` 同形的条目
* (`snapshotPiSessions` 用到的那四个字段 + created)。
*
* 只在这里做 I/O;调用方拿到的是纯数据。
*/
async function scan() {
const files = await collectFiles();
const alive = new Set(files);
// 会话文件被删掉之后缓存里的条目必须走,否则这张表就是下一个泄露源。
for (const key of [...cache.keys()]) {
if (!alive.has(key)) cache.delete(key);
}
const out = [];
// 并发 16:这些都是小 I/O(稳态下多数只有一次 stat),并发高一点能盖住
// 磁盘延迟;再高就只是给事件循环添堵。
for (let i = 0; i < files.length; i += 16) {
const batch = await Promise.all(
files.slice(i, i + 16).map((f) => scanOne(f).catch(() => null)),
);
for (const item of batch) if (item) out.push(item);
}
return out;
}
async function scanOne(file) {
let st;
try {
st = await stat(file);
} catch {
cache.delete(file);
return null;
}
const cached = cache.get(file);
// 文件变小 = 被重写/截断(pi 只 append,所以这不该发生 —— 但真发生时
// 缓存里的 size 会让我们从一个越界的偏移开始读)。整份重扫最保险。
const shrank = cached && st.size < cached.size;
if (!cached || shrank) {
const header = await readHeader(file);
if (!header) {
// 首行不是 header:不是会话文件。记一条空壳避免每拍都重读它。
cache.set(file, { size: st.size, id: '', cwd: '', created: new Date(0), name: undefined, hasName: false });
return null;
}
fullScans++;
tailBytes += st.size;
const { name, found } = await scanRangeForName(file, 0, st.size);
const entry = { size: st.size, id: header.id, cwd: header.cwd, created: header.created, name, hasName: found };
cache.set(file, entry);
return toInfo(file, entry, st);
}
// 空壳(已知不是会话文件):文件长了也不用管,它不会突然变成会话。
if (!cached.id) {
cached.size = st.size;
return null;
}
if (st.size > cached.size) {
tailScans++;
tailBytes += st.size - cached.size;
const { name, found } = await scanRangeForName(file, cached.size, st.size);
cached.size = st.size;
// 尾巴里没有 session_info 时**保留**旧 name。写成 `cached.name = name`
// 会让每次有新消息的会话都丢掉名字 —— 而没有 name 的会话不上报
// (S-1),于是活跃会话会从补全候选里消失。
if (found) {
cached.name = name;
cached.hasName = true;
}
}
return toInfo(file, cached, st);
}
function toInfo(file, entry, st) {
return {
path: file,
id: entry.id,
cwd: entry.cwd,
name: entry.name,
created: entry.created,
// 排序用「最近活跃」,mtime 就是它,且已经在手上(stat 已经做过了)
modified: st.mtime,
};
}
/** 观测用:稳态下 fullScans 应当不再增长,tailBytes 每拍只涨一点。 */
const stats = () => ({
tracked: cache.size,
fullScans,
tailScans,
tailBytes,
});
return { scan, stats };
}