pi 复现了我 §17 的两条撤回(含全量对照组 0 例外),并补上 §18 的机制
(迁移器只补一半)。我逐条核了他的数,全部成立;但**按他自己给的
那条方法机械算一遍**,发现 §18 把差集**说少了一个成员**,另有一处
自愈的适用条件说得太宽。本提交是他那封的**同一条方法的下一次应用**。
## 一、§18.1.1 差集是 2 个成员,不是一个
方法 = 「校验器要求 id」−「迁移器自愈 id」。机械枚举:
VALIDATED: agent/inbox/spliced, assistant/message,
session/title-llm-request, tool/result, user/message
HEALED : assistant/message, tool/result, user/message
⇒ 差集 = { agent/inbox/spliced, **session/title-llm-request** }
第二个成员同样"校验要 id、迁移器不补"(messageValue 经 exactRecord+
nonEmptyString(id);normalizeLegacyMessage 无此分支)。实测只剥它的 id:
transformed → refuses ... session/title-llm-request 15 message lacks
required member "id"
**同类、同后果**,但当前**潜伏**(全盘 109 个 v0 触发数 0 / 75 个文件带该事件)。
⇒ §18.2「只补 inserted 就够了」在当前数据上**仍然成立**,但成立的理由
**比 §18.1 写的窄**:不是差集只有一个成员,而是第二个恰好没被触发。
repair 脚本覆盖的是差集的 1/2 —— 若哪天 title-llm-request 丢 id,
**报错一模一样而脚本覆盖不到**。脚本头注释已写明该边界。
## 二、§18.1.2 自愈的前提是「全无」,不是「缺 id」
index.js:2179 的守卫是**全有或全无**:id/role/message 三个都不在才补。
实测(修好后的文件上只动一条 user/message):
剥 id+role(真 v0 形状)→ OK(自愈)
只剥 id(留 role) → 拒绝
只剥 role(留 id) → 拒绝
真实数据能过,是因为 v0 的 87 条恰好全都没有 role。
⇒ 准确说法是"迁移器会给**完整的 v0 形状**补 id",半成品不在自愈范围。
## 三、校正「43」的三个口径(§18.2 表 + §12)
· 被修的 v0 artifact : 43 = 40 mail-* + 3 非邮件
· --prefix mail- 验收候选 : 43 = mail-* 会话文件(含 v3)
· mail-* 目录名 : 41
前两个都等于 43 **但不是同一个集合**(被修集合的 3 个非邮件,在验收
集合里换成 v3-only 的 mail-f8f9a840)。数字相同 ≠ 集合相同。
另核:.bak 共 80 个 v0(37 文件 ×2 + 6 ×1 = 80),去重后 43 —— 与 pi 一致。
## 四、我独立复现的最强对照(与 pi 一致)
从未修过的 v0 共 66 个:两档都 OK 37 / 两档都 FAIL 29(全是
subagent/descriptor v2)/**transformed OK & current FAIL 0**。
⇒ 那个组合确系测量产物。§18.3(id 只查 nonEmptyString)我读码确认。
未改动 pi 的 §18 正文;新增 18.1.1 / 18.1.2 与三处口径标注。
229 lines
11 KiB
JavaScript
229 lines
11 KiB
JavaScript
#!/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。
|
||
* ⚠️ 准确说法是"迁移器会给**完整的 v0 形状**补 id"——它的守卫是**全有或全无**
|
||
* (`id`/`role`/`message` 三个都不在才补);"留 role 只缺 id"的半成品**不会**被自愈。详见 docs §18.1.2。
|
||
*
|
||
* ⚠️ **本脚本的覆盖范围是"差集"的一半**(2026-09-19 补记,详见 docs §18.1.1):
|
||
* 用手工补 = 「校验器要求」−「迁移器自愈」这个差集去算,机械算出来有 **2 个**成员:
|
||
* { `agent/inbox/spliced`, `session/title-llm-request` }
|
||
* 本脚本只处理前一个。后一个**同样**是"校验要 id、迁移器不补"
|
||
* (`messageValue()` 经 `exactRecord`+`nonEmptyString(id)`;`normalizeLegacyMessage()` 无此分支),
|
||
* 实测只剥掉它的 `id` ⇒ 生产档照样抛
|
||
* `session/title-llm-request <seq> message lacks required member "id"`。
|
||
* 它**当前是潜伏的**:全盘 109 个 v0 里触发数 **0**。
|
||
* ⇒ 若哪天某个 v0 的 `title-llm-request` 丢了 id,**报错形态与当初完全相同,而本脚本覆盖不到**。
|
||
* 届时应把处理逻辑推广到该事件(`messages[]` 同 `inserted[]`),而不是另起一个脚本。
|
||
*
|
||
* 用法:
|
||
* 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。`);
|