Files
MailUI4Agents/scripts/repair-v3-usermessage-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

144 lines
6.2 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
/**
* 修 v3 会话里 `user/message` 缺 `id`/`role` 的事件。
*
* 为什么需要它(与 v0 那个 repair 的区别):
* - `repair-legacy-spliced-ids.mjs` 修的是 **v0** 的 `agent/inbox/spliced.inserted[]`。
* - 修完之后,dsh 一旦 resume 该会话,就会把 v0 迁移成 **v3**(`session.v3.jsonl.zstd`)
* 并按当时的写入路径继续追加。写入路径会给迁移出来的历史消息补 id
* (`legacy-message:<sid>:<seq>`),但**新注入的**那条 `user/message`
* 仍然是 `{content, source}` —— 没有 id/role。
* - 生产读路径 `JsonlSessionPersistence.open()` 会走
* `validateStoredEvents → adoptSessionEvent → assertMessageEventShape`,
* 要求每条 message 事件都带非空 `id`;于是整个会话读不出来。
* 而 `readSession()` 因为 `SessionCorpus.load` 对 live 会话直接返回内存快照
* (不经磁盘校验),会**成功** —— 这就是两条路径结论相反的原因。
*
* 本脚本只做一件事:给 v3 里缺 `id`/`role` 的 `user/message` 补上,其余字节原样保留。
*
* ⚠️ 补进去的 id 是**伪造**的(原数据里没有)。它足以过生产档校验、让会话重新可读,
* 但不是"历史的真实 id"。治本仍在写入侧规范化。
*
* 用法:
* node scripts/repair-v3-usermessage-ids.mjs # dry-run(默认)
* node scripts/repair-v3-usermessage-ids.mjs --apply # 写盘(须先停 dsh.service)
*/
import { execFileSync } from "node:child_process";
import { readdirSync, statSync, copyFileSync, renameSync, writeFileSync, mkdirSync, rmSync, chmodSync, existsSync } from "node:fs";
import { join, basename } from "node:path";
import { tmpdir } from "node:os";
const dsRoot = process.env.DSH_INSTALL ?? "/usr/lib/node_modules/@deepseek-ai/dsh";
const { sessionFormatCatalog } = await import(join(dsRoot, "node_modules/@deepseek-ai/dsh-session-format-catalog/lib/index.js"));
/** v3 是当前版本,两级校验都应通过;只要 current 档过就算修好。 */
function v3Ok(header, events) {
for (const validation of ["transformed", "current"]) {
try {
const r = sessionFormatCatalog.createRestore(structuredClone(header), { recovery: "recoverable", validation });
for (const e of structuredClone(events)) r.decodeRow(e);
r.finish();
} catch (err) {
return { ok: false, msg: `[${validation}] ${err?.message}` };
}
}
return { ok: true };
}
const argv = process.argv.slice(2);
const APPLY = argv.includes("--apply");
const ROOT = "/root/.dsh/sessions";
function serviceActive() {
try {
return execFileSync("systemctl", ["is-active", "dsh.service"], { encoding: "utf8" }).trim() === "active";
} catch { return false; }
}
const readLines = (p) => execFileSync("zstd", ["-dc", p], { maxBuffer: 1 << 30 }).toString("utf8").split("\n").filter((l) => l.trim().length);
/** 生产同款帧布局:第 1 帧恰好一行 header,其后每 500 行一帧。 */
function writeArtifact(path, header, events) {
const tmp = join(tmpdir(), `dsh-v3fix-${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;
}
if (APPLY && serviceActive()) {
console.error("拒绝执行:dsh.service 正在运行。请先 `systemctl stop dsh.service` 再 --apply。");
process.exit(2);
}
const targets = [];
for (const store of readdirSync(ROOT)) {
const storePath = join(ROOT, store);
if (!statSync(storePath).isDirectory()) continue;
for (const id of readdirSync(storePath)) {
const p = join(storePath, id, "session.v3.jsonl.zstd");
if (existsSync(p)) targets.push({ id, path: p });
}
}
console.log(`模式: ${APPLY ? "APPLY(写盘)" : "dry-run(不写盘)"} v3 候选: ${targets.length}\n`);
let fixed = 0, clean = 0;
const details = [];
for (const { id, path } of targets) {
const lines = readLines(path);
const header = JSON.parse(lines[0]);
const events = lines.slice(1).map((l) => JSON.parse(l));
let patched = 0;
const out = events.map((e) => {
if (e.type !== "user/message") return e;
const d = e.data ?? {};
if (typeof d.id === "string" && d.id !== "" && d.role === "user") return e;
patched++;
return { ...e, data: { ...d, id: `recovered-usermessage-${e.seq}`, role: "user" } };
});
if (patched === 0) { clean++; continue; }
fixed++;
details.push([id, patched, events.length]);
console.log(` 需修复: ${id}(补 ${patched} 条 user/message,共 ${events.length} 事件)`);
if (APPLY) {
const backup = `${path}.bak-${new Date().toISOString().replace(/[:.]/g, "-")}`;
const mode = statSync(path).mode & 0o777;
copyFileSync(path, backup);
const staged = writeArtifact(path, header, out);
chmodSync(staged, mode);
// 写盘前自检:落盘字节必须两级校验都过,否则放弃(原文件未动)
try {
const l2 = readLines(staged);
const check = v3Ok(JSON.parse(l2[0]), l2.slice(1).map((x) => JSON.parse(x)));
if (!check.ok) throw new Error(check.msg);
} catch (err) {
rmSync(staged, { force: true });
console.error(` !! 自检失败,已放弃(原文件未动): ${String(err?.message).slice(0, 140)}`);
fixed--;
continue;
}
renameSync(staged, path);
console.log(` 已修复;备份 ${basename(backup)}`);
}
}
console.log(`\n=== 汇总 ===`);
console.log(` 已干净(无需处理): ${clean}`);
console.log(` 需修复: ${fixed}${APPLY ? "(已写入)" : "(dry-run:未写)"}`);