Files
MailUI4Agents/scripts/repair-legacy-spliced-ids.mjs
JianFeeeee 45577aa3b6 更正: 脚本注释里那条 别补 user/message 也是假象 —— 迁移器本来就会给它合成 id
与 §17 同源:后半句依据的是 校验污染入参后再校验 的假象。
深拷贝重测 40/40:补与不补都 strict-ok。
2026-09-19 12:16:08 +08:00

216 lines
9.8 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.

#!/usr/bin/env node
/**
* 修复 dsh 0.1.5 历史会话:`agent/inbox/spliced` 的 `inserted[]` 缺 `id`/`role`。
*
* 背景(实测,非推断):
* - `~/.dsh/sessions/**` 里 109 个会话是 v0(`session.jsonl.zstd`);新版 0.1.5-rc.2
* 读 v0 时要跑官方迁移链 v0→v1→v2→v3。
* - v0 解码器 `dsh-session-format-v0-to-v1` 的 `messageValue()` 要求 `inserted[]`
* 里每条消息都有 `id` 和 `role`;而**旧版写入时只写了 `{content, source}`**。
* ⇒ 迁移在 v0→v1 第一步就抛
* `agent/inbox/spliced <seq> inserted message lacks required member "id"`。
* - 于是桥读到的是 SessionQueryError,按“不在磁盘”处理 ⇒ 再 create 就撞
* `session "mail-…" already exists`。这就是“邮件通道全断”的根因。
*
* 本脚本做的事(**只改一处**):
* 给 `agent/inbox/spliced` 事件里缺 `id`/`role` 的 `inserted[]` 消息补上
* `id`(确定性、可复现)与 `role: "user"`,其余字节原样保留。
* 本脚本只补 `agent/inbox/spliced` 的 `inserted[]`。
* 关于 `user/message`(很多会话里它同样缺 id):
* 迁移器**本来就会**给它合成 id(对原始 v0 实测:迁出 386 条带 id / 0 条缺 id),
* 所以这里**不需要**补。早先注释说"补了会把能读的会话变成 seq gap"——
* **那是错的**(2026-09-19 更正):那是「校验器原地改写入参、又被拿去再校验」造成的假象;
* 深拷贝重测后补与不补都 strict-ok(40/40)。详见 docs §17。
*
* 用法:
* node scripts/repair-legacy-spliced-ids.mjs # 预演(默认,不写盘)
* node scripts/repair-legacy-spliced-ids.mjs --apply # 实际写入(须先停 dsh.service)
* node scripts/repair-legacy-spliced-ids.mjs --root <dir> --only mail-
*
* 安全性:
* - 默认 dry-run;只有显式 `--apply` 才写盘。
* - 写盘前逐会话用**生产同款选项**(recovery:"recoverable", validation:"transformed")
* 重跑迁移验证,验证不过就跳过、绝不写。
* - 原文件先复制为 `<file>.bak-<时间戳>`,再原子 rename 覆盖。
* - `--apply` 时若 dsh.service 在跑,直接拒绝(会话是单写者,必须离线)。
*/
import { execFileSync } from "node:child_process";
import { readdirSync, statSync, existsSync, copyFileSync, renameSync, writeFileSync, mkdirSync, rmSync, chmodSync } from "node:fs";
import { join, basename, dirname } from "node:path";
import { tmpdir } from "node:os";
const dsRoot = process.env.DSH_INSTALL ?? "/usr/lib/node_modules/@deepseek-ai/dsh";
const CATALOG = join(dsRoot, "node_modules/@deepseek-ai/dsh-session-format-catalog/lib/index.js");
const { sessionFormatCatalog } = await import(CATALOG);
// 生产读路径用的选项(dsh-session-persistence-jsonl: generationFormat)
const PROD = { recovery: "recoverable", validation: "transformed" };
const argv = process.argv.slice(2);
const APPLY = argv.includes("--apply");
const argOf = (flag, dflt) => {
const i = argv.indexOf(flag);
return i >= 0 && argv[i + 1] ? argv[i + 1] : dflt;
};
const ROOT = argOf("--root", "/root/.dsh/sessions");
const ONLY = argOf("--only", "");
function serviceActive() {
try {
return execFileSync("systemctl", ["is-active", "dsh.service"], { encoding: "utf8" }).trim() === "active";
} catch {
return false; // is-active 非 active 时退出码非 0
}
}
function readLines(path) {
const raw = execFileSync("zstd", ["-dc", path], { maxBuffer: 1 << 30 }).toString("utf8");
return raw.split("\n").filter((l) => l.trim().length > 0);
}
/** 验证器会原地改写入参,凡是「验证 + 写盘」共用的数据都要先深拷贝。 */
const clone = (o) => structuredClone(o);
function migrateOk(header, events) {
try {
// ⚠️ 验证器会**原地改写**入参:实测 `reasoning-chunks` 等事件在 decodeRow 后
// 会被写回补全的 `dt` 数组(179 事件里有 5 个被改)。若拿同一批对象既验证
// 又写盘,写出去的就是**被验证器污染过的数据**,重读时报
// `released Session row N has seq gap` —— 这会凭空把「可修复」变成「不可修复」。
// 因此验证一律在深拷贝上进行,出参保持原始字节语义。
const r = sessionFormatCatalog.createRestore(clone(header), PROD);
for (const e of clone(events)) r.decodeRow(e);
r.finish();
return { ok: true };
} catch (err) {
return { ok: false, msg: err?.message ?? String(err) };
}
}
/** 只给 agent/inbox/spliced 的 inserted[] 补 id/role;其余行原样返回。 */
function patchLines(path, lines) {
const header = JSON.parse(lines[0]);
const events = lines.slice(1).map((l) => JSON.parse(l));
const before = migrateOk(header, events);
if (before.ok) return { status: "already-readable", header, events, patched: 0 };
const splicedRe = /"agent\/inbox\/spliced"/;
let patched = 0;
const outEvents = events.map((e) => {
if (e.type !== "agent/inbox/spliced") return e;
const inserted = (e.data.inserted ?? []).map((m) => {
if ("id" in m && "role" in m) return m;
patched++;
// 确定性 id:同一事件同一位置永远得到同一个 id,重跑可复现
return { id: `recovered-splice-${e.seq}-${patched}`, role: "user", ...m };
});
return { ...e, data: { ...e.data, inserted } };
});
if (patched === 0) return { status: "unrepairable", header, events, before: before.msg, patched: 0 };
const after = migrateOk(header, outEvents);
return { status: after.ok ? "repairable" : "unrepairable", header, events, outEvents, before: before.msg, after: after.msg, patched };
}
/** 按生产布局重写:第 1 帧恰好一行 header,其后每 500 行一帧。 */
function writeArtifact(path, header, events) {
const tmp = join(tmpdir(), `dsh-repair-${process.pid}-${Date.now()}`);
mkdirSync(tmp, { recursive: true });
const parts = [];
const h = join(tmp, "h.jsonl");
writeFileSync(h, JSON.stringify(header) + "\n");
execFileSync("zstd", ["-q", "-f", h, "-o", join(tmp, "f0.zst")]);
parts.push(join(tmp, "f0.zst"));
const BATCH = 500;
for (let i = 0, k = 0; i < events.length; i += BATCH, k++) {
const b = join(tmp, `b${k}.jsonl`);
writeFileSync(b, events.slice(i, i + BATCH).map((o) => JSON.stringify(o)).join("\n") + "\n");
execFileSync("zstd", ["-q", "-f", b, "-o", join(tmp, `f${k + 1}.zst`)]);
parts.push(join(tmp, `f${k + 1}.zst`));
}
const staged = `${path}.repair-staged`;
execFileSync("bash", ["-c", `cat ${parts.map((p) => JSON.stringify(p)).join(" ")} > ${JSON.stringify(staged)}`]);
rmSync(tmp, { recursive: true, force: true });
return staged;
}
function* walk(dir) {
for (const name of readdirSync(dir)) {
const p = join(dir, name);
if (statSync(p).isDirectory()) yield* walk(p);
else if (name === "session.jsonl.zstd") yield p;
}
}
// ---- main ----
if (APPLY && serviceActive()) {
console.error("拒绝执行:dsh.service 正在运行。会话日志是单写者,请先 `systemctl stop dsh.service` 再 --apply。");
process.exit(2);
}
const targets = [...walk(ROOT)].filter((p) => (ONLY ? p.includes(ONLY) : true));
console.log(`root: ${ROOT} 模式: ${APPLY ? "APPLY(写盘)" : "dry-run(不写盘)"} 候选文件: ${targets.length}\n`);
const tally = { "already-readable": 0, repairable: 0, unrepairable: 0 };
const unrepairable = [];
for (const path of targets) {
let lines;
try {
lines = readLines(path);
} catch (err) {
console.log(` 跳过(解压失败): ${path}: ${String(err?.message).slice(0, 80)}`);
continue;
}
const header = JSON.parse(lines[0]);
if (header.version !== 0) { tally["already-readable"]++; continue; }
const info = patchLines(path, lines);
const id = header.id;
if (info.status === "already-readable") { tally["already-readable"]++; continue; }
if (info.status === "unrepairable") {
tally.unrepairable++;
unrepairable.push([id, info.after ?? info.before]);
continue;
}
tally.repairable++;
console.log(` 可修复: ${id} (补 ${info.patched} 条 inserted 消息)`);
if (APPLY) {
// 逐会话隔离失败:单个会话写不动(EACCES/只读/并发)不该中止整轮,
// 否则前面已修的与后面待修的都会因为"一次抛错"被跳过。
let staged;
try {
const backup = `${path}.bak-${new Date().toISOString().replace(/[:.]/g, "-")}`;
const mode = statSync(path).mode & 0o777;
copyFileSync(path, backup);
staged = writeArtifact(path, header, info.outEvents);
chmodSync(staged, mode);
// 写盘后立刻用生产读路径自检
const check = (() => {
try {
const l2 = readLines(staged);
return migrateOk(JSON.parse(l2[0]), l2.slice(1).map((x) => JSON.parse(x)));
} catch (err) { return { ok: false, msg: err?.message ?? String(err) }; }
})();
if (!check.ok) throw new Error(`写盘前自检失败: ${check.msg}`);
renameSync(staged, path);
console.log(` 已修复;备份 ${basename(backup)}`);
} catch (err) {
if (staged) rmSync(staged, { force: true });
console.error(` !! 跳过(原文件未动): ${String(err?.message).slice(0, 140)}`);
tally.repairable--; tally.unrepairable++;
unrepairable.push([id, `apply 失败: ${err?.message}`]);
continue;
}
}
}
console.log(`\n=== 汇总 ===`);
console.log(` 已可读(无需处理): ${tally["already-readable"]}`);
console.log(` 可修复: ${tally.repairable}${APPLY ? "(已写入)" : "(dry-run:未写)"}`);
console.log(` 仍不可修复: ${tally.unrepairable}`);
for (const [id, msg] of unrepairable.slice(0, 20)) console.log(` - ${id}: ${String(msg).slice(0, 120)}`);
if (!APPLY && tally.repairable > 0) console.log(`\n确认无误后:先停 dsh.service,再跑 --apply。`);