Files
MailUI4Agents/scripts/repair-legacy-spliced-ids.mjs
JianFeeeee b4a8f74ae5 修复: dsh 邮件通道全断的**两侧**根因(桥侧不产 message id 是真正在写的那一处)
现象:dsh 的邮件通道全断。老会话读不出来 ⇒ 桥报 SessionQueryError ⇒ 按"不在磁盘"
处理 ⇒ 再 create 撞 `already exists`。修好读路径之后又立刻暴露下一层
`message "undefined" is already pending`。

根因一(历史数据,dsh 侧):v0 会话的 `agent/inbox/spliced.inserted[]` 缺 `id`/`role`,
v0→v1 迁移第一步就拒绝。40 个真 mail-* 会话全部命中。

根因二(**仍在写**,本仓侧):`plugins/dsh-mail-bridge/lib/message.js` 的
`userMessage()` 只产出 `{content, source}`。DSH 0.1.5 的 inbox 按 `message.id` 去重
(`dsh-agent-loop` 的投影 apply() 与 mutate() 各维护一个 Set),id 全是 undefined
⇒ **第二条消息必挂**。日志里最早的同类记录在 2026-09-07,累计 50+ 次。
官方形状在 `@deepseek-ai/dsh-llm` 的 `createMessage()`({id, role, content, source}),
同一份 dsh 里其它插件都用官方的 createUserMessage(),只有这个桥手搓。
以前没炸是因为读路径先坏,根本走不到 followup。

本次改动
- message.js/.d.ts: userMessage() 补 id: randomUUID() 与 role:'user'
- test/message.test.mjs: 钉住「id 非空」「两条消息 id 必须不同」,用官方 inbox
  去重逻辑逐字复刻验证(修复前 message "undefined" is already pending,修复后 20 封全唯一)
- scripts/: repair-legacy-spliced-ids.mjs(v0,默认 dry-run)、
  repair-v3-usermessage-ids.mjs(v3)、verify-mail-sessions-readable.mjs
  (走生产真读路径 JsonlSessionPersistence.open,而非解码器口径)、两个 apply driver
- docs/DSH-0.1.5-MAIL-CHANNEL-ROOTCAUSE.md: 补执行结果与两处新事实

执行与验收(详见文档 §9-§15)
- v0 修 40 个、v3 修 2 个;逐文件解压后与备份 `cmp` **逐字节相等**,事件数 40/40 一致,
  零丢失(25.2MB→12.5MB 是单帧改 500 行/帧的重压缩,不是丢数据)
- 真 mail-* 会话最终 **41/41 可读**
- journal 里同一会话从 `already exists` 变为 `resume 续谈`,且持续增长
  (22647→22685 事件),最新 user/message 带真实 UUID;修复上线后 already pending 计数为 0
- 已在生产部署(deploy/redeploy-plugin.sh dsh,快照+原子软链+重启+后置验证全绿)

两个必须记住的坑
1. **校验与落盘不能共用同一批对象**:createRestore().decodeRow() 会原地改写入参
   (补全 dt 数组),污染后写出去会报 `released Session row N has seq gap`。
   这曾让 dry-run 说"40 个可修"、apply 只说"3 个"。
2. **判定磁盘健康只认 open()**:readSession() 走 SessionCorpus.load,命中有 live 会话时
   直接返回内存快照、不校验磁盘;open() 才走 validateStoredEvents。两条路径结论相反
   是设计使然,不是矛盾。
2026-09-19 12:03:34 +08:00

213 lines
9.4 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"`,其余字节原样保留。
*
* 明确**不要**顺手给 `user/message` 也补 id(很多会话里它同样缺 id):
* 实测那样做会把本来能读的会话变成 `seq gap` 而读不了。
*
* 用法:
* 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。`);