#!/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 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 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 --only mail- * * 安全性: * - 默认 dry-run;只有显式 `--apply` 才写盘。 * - 写盘前逐会话用**生产同款选项**(recovery:"recoverable", validation:"transformed") * 重跑迁移验证,验证不过就跳过、绝不写。 * - 原文件先复制为 `.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。`);