From 0446a9e8f050c6e392dc2c20a7ba8040f02eddab Mon Sep 17 00:00:00 2001 From: zizi Date: Thu, 23 Jul 2026 11:49:26 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D:=20=E6=94=B6=E7=B4=A7?= =?UTF-8?q?=E5=8D=87=E6=A0=BC=E5=85=B3=E7=B3=BB=E9=87=8D=E6=95=B4=E4=B8=8E?= =?UTF-8?q?=E8=A1=A5=E5=81=BF=E8=BE=B9=E7=95=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .claude/skills/parse-book/SKILL.md | 4 +- .../parse-book/scripts/parse_upgrade.py | 188 ++++++++-- .../scripts/test_parse_upgrade_offline.py | 353 +++++++++++++++++- docs/2026-07-16-升格卡改造设计.md | 7 +- 4 files changed, 511 insertions(+), 41 deletions(-) diff --git a/.claude/skills/parse-book/SKILL.md b/.claude/skills/parse-book/SKILL.md index 3b3a3a6..f53fe50 100644 --- a/.claude/skills/parse-book/SKILL.md +++ b/.claude/skills/parse-book/SKILL.md @@ -58,7 +58,9 @@ disable-model-invocation: true **窗行陷阱(放量首日实测)**:`--window` 参数变化后重切,旧窗行会按 from_order 占位,新的大窗被「已有大纲跳过」→ 中间章域永远漏出卡(验收期 1–3 章小窗占住 from_order=1,放量 1–34 章大窗被跳过)。**换窗参数重切前必须先删该书全部窗行**(窗行是可再生中间产物;卡挂「窗起」,cards 重出时按窗软删重出)。 -**作品面升格执行器 `scripts/parse_upgrade.py`(命令 `windows`/`run`/`status`)**:正文按窗抽取 `upgrade_book` 实体卡。每窗采用“两阶段短事务”:模型/嵌入调用期间不持有业务连接;实体写入和关系写入各自提交 `processing` marker;最终嵌入必须完整成功,才与 `done` 在同一短事务提交。窗口输入摘要绑定作品标题、章/块 ID、标题、状态、revision、类型、正文、active schema 合同和当前脚本 SHA。启动时会恢复可精确撤销的 `processing`;本次窗口第二次失败立即以非零退出并停止当前作品,保留 `retryable-clean` 断点,人工再次运行才继续。 +**作品面升格执行器 `scripts/parse_upgrade.py`(命令 `windows`/`run`/`status`)**:正文按窗抽取 `upgrade_book` 实体卡。每窗采用“两阶段短事务”:模型/嵌入调用期间不持有业务连接;实体写入和关系写入各自提交 `processing` marker;最终嵌入必须完整成功,才与 `done` 在同一短事务提交。新卡 `INSERT` 必须返回数据库实际 `revision`,同窗后续更新的 CAS 绑定该值,禁止假定初始版本。关系首次同 pair 冲突只允许一次模型 repair;repair 仍拆成多条时,只机械合并关系类型、本窗演变和无字段冲突的其他维度,非空未知顶层键或同字段异值继续失败关闭,不增加第三次模型调用。窗口输入摘要绑定作品标题、章/块 ID、标题、状态、revision、类型、正文、active schema 合同和当前脚本 SHA。启动时会恢复可精确撤销的 `processing`;本次窗口第二次失败立即以非零退出并停止当前作品,保留 `retryable-clean` 断点,人工再次运行才继续。 + +普通无 marker clean 的产物检查只豁免满足完整墓碑条件且 `updater` 为 `upgrade-reset` 或 `upgrade-undo` 的 alias/new_card/embedding;presence、card_state、audit 仍有行即阻断。下述 legacy 恢复只允许 `upgrade-reset`,两种上下文不得混用。 `recover-legacy-failed --preview/--execute` 是 fence 引入前 failed 窗的唯一恢复入口:preview 只读输出窗口、六类本窗产物计数、processing 数和确认 SHA;execute 需确认旧进程已结束、同书锁内重算并 exact 匹配,且无 processing、**阻断产物计数为 0**,才保持 `failed` 并标为 `retryable-clean`。计数只豁免能机械证明来源的全书 reset 墓碑:draft 必须同时 `deleted=true`、`updater='upgrade-reset'`、`status='pending'`;embedding 还必须 `entity_id IS NULL` 且自身 `deleted=true`;alias 必须 `deleted=true` 且 `updater='upgrade-reset'`。presence 没有 reset 来源标记,任何本窗行(包括软删墓碑)都阻断;其他软删、活跃行、已有 entity owner,以及任何 card_state/audit 行也都阻断。完整 `stateSha` 仍覆盖所有 deleted 行,二者不能混淆。不调用模型/嵌入。`--redo-window` 已禁用,历史修正走人工 backup/reset/rebuild。 diff --git a/.claude/skills/parse-book/scripts/parse_upgrade.py b/.claude/skills/parse-book/scripts/parse_upgrade.py index c7c9bd4..2beaf72 100644 --- a/.claude/skills/parse-book/scripts/parse_upgrade.py +++ b/.claude/skills/parse-book/scripts/parse_upgrade.py @@ -51,6 +51,9 @@ WIN_MAX_CHAPS = 12 # 单窗章数上限 UPDATE_BATCH = 6 # 卡更新每批张数上限(防长清单丢字段) SOURCE_TYPE = "upgrade_book" # 独立来源标记:与范式卡 parse_book 隔离判重/检索/确认 +# 墓碑 updater 只允许来自明确的恢复/撤销路径;上下文再从中选择自己的子集。 +TOMBSTONE_UPDATER_ALLOWLIST = frozenset(("upgrade-reset", "upgrade-undo")) + # 作品面七型(六实体型+关系型;大纲卡全书收尾单独做,不在窗循环内) ENTITY_TYPES = ("character", "location", "item", "faction", "power_system", "event") RELATION_TYPE = "character_relation" @@ -1064,11 +1067,9 @@ def _normalize_relation_output(relations, valid_draft_ids): for relation in relations or []: if not isinstance(relation, dict): continue - first, second = relation.get("甲方"), relation.get("乙方") - # 保持既有语义:不在本窗核心角色内的锚和自关系直接忽略。 - if first not in valid_draft_ids or second not in valid_draft_ids or first == second: + pair = _normalized_relation_pair(relation, valid_draft_ids) + if pair is None: continue - pair = tuple(sorted((first, second))) item = deepcopy(relation) item["甲方"], item["乙方"] = pair if pair not in relation_by_pair: @@ -1082,6 +1083,122 @@ def _normalize_relation_output(relations, valid_draft_ids): return normalized +def _normalized_relation_pair(relation, valid_draft_ids): + """返回合法关系的无序 pair;窗外实体和自关系沿用原规则直接过滤。""" + + first, second = relation.get("甲方"), relation.get("乙方") + if first not in valid_draft_ids or second not in valid_draft_ids or first == second: + return None + return tuple(sorted((first, second))) + + +def _relation_value_parts(value): + """把关系维度转为可稳定去重的非空文本片段,不丢弃非空列表成员。""" + + if not _relation_value_is_nonempty(value): + return [] + if isinstance(value, (list, tuple)): + parts = [] + for item in value: + parts.extend(_relation_value_parts(item)) + return parts + text = str(value).strip() + return [text] if text else [] + + +def _relation_value_is_nonempty(value): + """识别 JSON 维度是否携带信息;空字符串、空容器和 null 可安全忽略。""" + + if value is None: + return False + if isinstance(value, str): + return bool(value.strip()) + if isinstance(value, (list, tuple, dict, set)): + return bool(value) + return True + + +def _merge_repaired_relation_output(relations, valid_draft_ids): + """只供第二次 repair 输出使用:按无序 pair 安全合并多维关系,不覆盖冲突值。""" + + standard_keys = frozenset(("甲方", "乙方", "关系类型", "本窗演变", "其他字段")) + states = {} + for relation in relations or []: + if not isinstance(relation, dict): + continue + pair = _normalized_relation_pair(relation, valid_draft_ids) + if pair is None: + continue + unknown_keys = [ + str(key) + for key, value in relation.items() + if key not in standard_keys and _relation_value_is_nonempty(value) + ] + if unknown_keys: + raise RelationOutputConflict( + f"关系 repair 含非空未知顶层键:甲方draft={pair[0]}," + f"乙方draft={pair[1]},keys={sorted(unknown_keys)}" + ) + if pair not in states: + states[pair] = { + "关系类型": [], + "本窗演变": [], + "其他字段": {}, + } + state = states[pair] + for field_name in ("关系类型", "本窗演变"): + for part in _relation_value_parts(relation.get(field_name)): + if part not in state[field_name]: + state[field_name].append(part) + + other_fields = relation.get("其他字段") + if other_fields is None: + other_fields = {} + if not isinstance(other_fields, dict): + raise RelationOutputConflict( + f"关系 repair 的其他字段不是对象:甲方draft={pair[0]},乙方draft={pair[1]}" + ) + for field_name, value in other_fields.items(): + if not _relation_value_parts(value): + continue + if field_name not in state["其他字段"]: + state["其他字段"][field_name] = deepcopy(value) + continue + if state["其他字段"][field_name] != value: + raise RelationOutputConflict( + f"关系 repair 同一字段值冲突:甲方draft={pair[0]}," + f"乙方draft={pair[1]},field={field_name}" + ) + + # 模型返回顺序不属于业务语义;输出前统一固定 pair、片段和字段键的顺序。 + return [ + { + "甲方": pair[0], + "乙方": pair[1], + "关系类型": " / ".join(sorted(states[pair]["关系类型"])), + "本窗演变": ";".join(sorted(states[pair]["本窗演变"])), + "其他字段": { + field_name: deepcopy(states[pair]["其他字段"][field_name]) + for field_name in sorted(states[pair]["其他字段"]) + }, + } + for pair in sorted(states) + ] + + +def _relation_pair_counts(relations, valid_draft_ids): + """统计 repair 输出中每个合法 pair 的原始条数,供二次合并日志使用。""" + + counts = {} + for relation in relations or []: + if not isinstance(relation, dict): + continue + pair = _normalized_relation_pair(relation, valid_draft_ids) + if pair is not None: + counts[pair] = counts.get(pair, 0) + 1 + return counts + + def _insert_relation_card(conn, work_id, win_no, payload): """新建关系卡并返回 draft id;逐字段审计和 state 任一步失败都会让窗事务回滚。""" @@ -1630,13 +1747,14 @@ def new_card( # 先原子占用全部 alias;任一不同 canonical 冲突都会让窗事务在写 payload 前失败关闭。 for alias in payload["别名"]: _claim_alias(conn, work_id, payload["名称"], alias, win_no, "init") - did = conn.execute( + did, initial_revision = conn.execute( """INSERT INTO muse_knowledge_draft (work_id, draft_type, draft_payload, status, source_type, source_id, creator, updater, tenant_id) - VALUES (%s,'entity',%s,'pending',%s,%s,'upgrade','upgrade',%s) RETURNING id""", + VALUES (%s,'entity',%s,'pending',%s,%s,'upgrade','upgrade',%s) + RETURNING id, revision""", (work_id, json.dumps(payload, ensure_ascii=False), SOURCE_TYPE, work_id, - TENANT)).fetchone()[0] + TENANT)).fetchone() for k, v in payload["字段"].items(): conn.execute( """INSERT INTO example_upgrade_audit @@ -1647,7 +1765,7 @@ def new_card( """INSERT INTO example_upgrade_card_state (draft_id, work_id, watermark_window, tenant_id) VALUES (%s,%s,%s,%s) ON CONFLICT (draft_id) DO NOTHING""", (did, work_id, win_no, TENANT)) - return did, payload["名称"], tuple(payload["别名"]) + return did, payload["名称"], tuple(payload["别名"]), initial_revision def undo_window(conn, work_id, win_no): @@ -2339,15 +2457,16 @@ def _apply_entity_stage(conn, work_id, win_no, plan, updates, snapshots, expecte existing_refs.update(ref for ref in plan["appearances"] if ref > 0) for ref in sorted(existing_refs): _lock_upgrade_draft(conn, ref, expected_revisions[ref]) - ref_map, new_ids = {}, set() + ref_map, new_ids, initial_revision_by_ref = {}, set(), {} for action in plan["actions"]: if action[0] == "new": _, ref, ent, history, chains = action - did, _, _ = new_card( + did, _, _, initial_revision = new_card( conn, work_id, win_no, ent, milestone_types, chapter_texts=chapter_texts, known_chapters=history, ) ref_map[ref], new_ids = did, new_ids | {did} + initial_revision_by_ref[ref] = initial_revision for other, relation, name in chains: conn.execute( """INSERT INTO example_upgrade_audit @@ -2373,7 +2492,7 @@ def _apply_entity_stage(conn, work_id, win_no, plan, updates, snapshots, expecte touched = set(new_ids) for update in updates: ref, actual = update["draft_id"], ref_map.get(update["draft_id"], update["draft_id"]) - revision = 0 if ref < 0 else snapshots[ref][1] + revision = initial_revision_by_ref[ref] if ref < 0 else snapshots[ref][1] merge_card( conn, actual, win_no, update["变更字段"], update["别名新增"], valid_keys=update["valid_keys"], chapter_texts=chapter_texts, @@ -2422,7 +2541,15 @@ def _compute_relations(contracts, title, a, b, text, characters, relation_rows, relation_repair_prompt(contracts, title, a, b, text, characters, relations, raw), ("关系",), ) - return _normalize_relation_output(repaired.get("关系") or [], valid_ids) + repaired_raw = repaired.get("关系") or [] + merged = _merge_repaired_relation_output(repaired_raw, valid_ids) + for pair, count in sorted(_relation_pair_counts(repaired_raw, valid_ids).items()): + if count > 1: + click.echo( + f"关系 repair 二次机械合并:pair={pair},合并数={count - 1}", + err=True, + ) + return merged def _apply_relation_stage(conn, work_id, win_no, relation_items, characters, relation_rows, @@ -2497,8 +2624,21 @@ def _finalize_window(work_id, win_no, input_sha, relation_state_sha, prepared): return embedded -def _window_artifact_counts(conn, work_id, win_no): - """统计阻断产物;仅机械证明的 upgrade-reset 墓碑允许豁免。""" +def _window_artifact_counts( + conn, + work_id, + win_no, + allowed_tombstone_updaters=("upgrade-reset",), +): + """统计阻断产物;墓碑豁免 updater 必须由调用上下文明确传入。""" + + if isinstance(allowed_tombstone_updaters, (str, bytes)): + raise ValueError("allowed_tombstone_updaters 必须是 updater 序列") + allowed_tombstone_updaters = tuple(dict.fromkeys(allowed_tombstone_updaters)) + unexpected = set(allowed_tombstone_updaters) - TOMBSTONE_UPDATER_ALLOWLIST + if unexpected: + raise ValueError(f"不允许的墓碑 updater:{sorted(unexpected)}") + allowed_updater_array = list(allowed_tombstone_updaters) values = {} statements = { @@ -2507,15 +2647,15 @@ def _window_artifact_counts(conn, work_id, win_no): (TENANT, work_id, win_no)), "aliases": ("SELECT count(*) FROM example_upgrade_alias a WHERE a.tenant_id=%s " "AND a.work_id=%s AND a.evidence_window=%s " - "AND NOT (a.deleted=TRUE AND a.updater='upgrade-reset')", - (TENANT, work_id, win_no)), + "AND NOT (a.deleted=TRUE AND a.updater = ANY(%s))", + (TENANT, work_id, win_no, allowed_updater_array)), "presence": ("SELECT count(*) FROM example_upgrade_presence p WHERE p.tenant_id=%s " "AND p.work_id=%s AND p.window_no=%s", (TENANT, work_id, win_no)), "new_cards": ("SELECT count(*) FROM muse_knowledge_draft d WHERE d.tenant_id=%s " "AND d.work_id=%s AND d.source_type=%s AND d.draft_payload->>'来源'=%s " - "AND NOT (d.deleted=TRUE AND d.updater='upgrade-reset' AND d.status='pending')", - (TENANT, work_id, SOURCE_TYPE, f"升格@窗{win_no}")), + "AND NOT (d.deleted=TRUE AND d.updater = ANY(%s) AND d.status='pending')", + (TENANT, work_id, SOURCE_TYPE, f"升格@窗{win_no}", allowed_updater_array)), "card_state": ("SELECT count(*) FROM example_upgrade_card_state s " "JOIN muse_knowledge_draft d ON d.id=s.draft_id " "WHERE s.tenant_id=%s AND d.tenant_id=%s AND d.work_id=%s " @@ -2526,8 +2666,9 @@ def _window_artifact_counts(conn, work_id, win_no): "AND d.work_id=%s AND d.source_type=%s " "AND d.draft_payload->>'来源'=%s " "AND NOT (e.entity_id IS NULL AND e.deleted=TRUE " - "AND d.deleted=TRUE AND d.updater='upgrade-reset' AND d.status='pending')", - (TENANT, TENANT, work_id, SOURCE_TYPE, f"升格@窗{win_no}")), + "AND d.deleted=TRUE AND d.updater = ANY(%s) AND d.status='pending')", + (TENANT, TENANT, work_id, SOURCE_TYPE, f"升格@窗{win_no}", + allowed_updater_array)), } for name, (sql, params) in statements.items(): values[name] = conn.execute(sql, params).fetchone()[0] @@ -2773,7 +2914,12 @@ def _compensate_window(work_id, win_no, error): raise CompensationFenceConflict("补偿失败 marker CAS 冲突") from exc conn.commit() raise CompensationFenceConflict(message) from exc - elif any(_window_artifact_counts(conn, work_id, win_no).values()): + elif any(_window_artifact_counts( + conn, + work_id, + win_no, + ("upgrade-reset", "upgrade-undo"), + ).values()): message = f"{COMPENSATION_FAILED_PREFIX} 实体提交前无 marker,但本窗产物非零"[:500] changed = conn.execute( """UPDATE example_upgrade_window SET status='failed', error_message=%s, updater='upgrade' diff --git a/.claude/skills/parse-book/scripts/test_parse_upgrade_offline.py b/.claude/skills/parse-book/scripts/test_parse_upgrade_offline.py index 3584976..9a9feb7 100644 --- a/.claude/skills/parse-book/scripts/test_parse_upgrade_offline.py +++ b/.claude/skills/parse-book/scripts/test_parse_upgrade_offline.py @@ -16,10 +16,12 @@ import json import hashlib import inspect +from io import StringIO import os import pathlib import sys from contextlib import nullcontext +from contextlib import redirect_stderr from copy import deepcopy from unittest.mock import patch @@ -334,7 +336,7 @@ class _CardConn: def __init__(self, payload=None): self.payload = payload self.audits = [] - self.revision = 0 + self.revision = 1 def execute(self, query, params): """按 SQL 用途返回最小结果,或捕获立卡、更新后的 payload。""" @@ -355,7 +357,8 @@ class _CardConn: return _CardResult(rows=[(101, self.payload)]) if "INSERT INTO muse_knowledge_draft" in query: self.payload = json.loads(params[1]) - return _CardResult((101,)) + return _CardResult((101, self.revision) + if "RETURNING id, revision" in normalized else (101,)) if normalized.startswith("INSERT INTO example_upgrade_alias"): return _CardResult((params[1],)) if "INSERT INTO example_upgrade_audit" in query and len(params) >= 4 \ @@ -878,7 +881,13 @@ class _RunConn: draft_id = self.state.setdefault("next_draft_id", 101) self.state["next_draft_id"] = draft_id + 1 self.state.setdefault("payloads", {})[draft_id] = json.loads(params[1]) - return _CardResult((draft_id,)) + initial_revision = self.state.get("new_card_initial_revision", 1) + self.state.setdefault("revisions", {})[draft_id] = initial_revision + self.state.setdefault("inserted_draft_revisions", {})[draft_id] = initial_revision + return _CardResult( + (draft_id, initial_revision) + if "RETURNING id, revision" in normalized else (draft_id,) + ) if normalized.startswith("UPDATE muse_knowledge_draft SET draft_payload=%s"): self.state.setdefault("draft_update_sql", []).append(normalized) if "RETURNING revision" in normalized: @@ -1042,6 +1051,67 @@ def test_run_new_card_registers_normalized_name_and_aliases_same_window(): {params[1] for params in candidate_rows} == {"白色游魂"}, detail=str(candidate_rows)) +def test_new_card_same_window_update_uses_database_initial_revision(): + """新卡同窗继续更新时,CAS 必须从 INSERT 返回的数据库初始 revision=1 起步。""" + + state = { + "payloads": {}, + "window_status": "pending", + "new_card_initial_revision": 1, + } + conn = _RunConn(state) + entity = { + "型": "item", + "名称": "星钥", + "别名": [], + "一句话摘要": "沉睡中的钥匙", + "字段": {}, + "出场章": [70, 71], + } + plan = { + "actions": [("new", -1, entity, set(), [])], + "to_update": {-1: ["同窗补充状态"]}, + "appearances": {}, + "aliases": {}, + } + updates = [{ + "draft_id": -1, + "变更字段": {"一句话摘要": "星钥已激活"}, + "别名新增": [], + "valid_keys": {"一句话摘要"}, + }] + chapter_texts = {70: "星钥被发现。", 71: "星钥完成激活。"} + input_sha = pu._window_input_sha(*pu._capture_window_input(conn, 8, 7)) + + with patch.object(pu, "merge_card", wraps=pu.merge_card) as merge_mock: + ref_map, _, _, touched, _ = pu._apply_entity_stage( + conn, + 8, + 7, + plan, + updates, + {-1: (pu._build_new_card_payload(8, 7, entity), 0)}, + {}, + set(), + chapter_texts, + input_sha, + 91, + ) + + draft_id = ref_map[-1] + check("new-card-revision-INSERT取得数据库初始值", + state["inserted_draft_revisions"][draft_id] == 1, detail=str(state)) + check("new-card-revision-同窗merge使用真实初始revision", + merge_mock.call_args.kwargs["expected_revision"] == 1, + detail=str(merge_mock.call_args)) + check("new-card-revision-CAS从1推进到2并完成实体阶段", + state["revisions"][draft_id] == 2 + and state["payloads"][draft_id]["一句话摘要"] == "星钥已激活" + and touched == {draft_id} + and state["window_status"] == "processing", + detail=str(state)) + + def _run_duplicate_entity_update_case(updates, *, return_error=False): """执行同一实体的重复更新输出,并返回写入口调用和窗失败证据。""" @@ -1403,7 +1473,7 @@ def _run_relation_repair_case(first_relations, repaired_relations, *, return_err def test_run_relation_conflict_repair_merges_to_one_card(): - """首轮同一无序人物对冲突时,关系专用重整合并后只写一张卡并完成整窗。""" + """repair 仍拆成同一无序人物对的多维记录时,机械合并后只写一张卡。""" first = [ {"甲方": 11, "乙方": 12, "关系类型": "盟友", "本窗演变": "并肩作战", @@ -1411,12 +1481,23 @@ def test_run_relation_conflict_repair_merges_to_one_card(): {"甲方": 12, "乙方": 11, "关系类型": "上下级", "本窗演变": "接受指挥", "其他字段": {"权力结构": "乙指挥甲"}}, ] - repaired = [{ - "甲方": 11, "乙方": 12, "关系类型": "盟友兼上下级", - "本窗演变": "并肩作战并接受指挥", - "其他字段": {"当前状态": "互信", "权力结构": "乙指挥甲"}, - }] - db, prompts, update_calls, insert_calls, undo_calls = _run_relation_repair_case(first, repaired) + repaired = [ + { + "甲方": 11, "乙方": 12, "关系类型": "盟友", + "本窗演变": "并肩作战", + "其他字段": {"当前状态": "互信"}, + }, + { + "甲方": 12, "乙方": 11, "关系类型": "上下级", + "本窗演变": "接受指挥", + "其他字段": {"权力结构": "乙指挥甲"}, + }, + ] + with StringIO() as stderr, redirect_stderr(stderr): + db, prompts, update_calls, insert_calls, undo_calls = _run_relation_repair_case( + first, repaired + ) + repair_log = stderr.getvalue() relation_cards = [ row for row in db.drafts.values() if row["payload"].get("type") == pu.RELATION_TYPE @@ -1431,12 +1512,143 @@ def test_run_relation_conflict_repair_merges_to_one_card(): check("relation-repair-合并后只写一张关系卡", update_calls == 0 and insert_calls == 1 and len(relation_cards) == 1, detail=f"update={update_calls}, insert={insert_calls}, cards={relation_cards}") + relation_payload = relation_cards[0]["payload"] + check("relation-repair-关系类型稳定合并", + relation_payload["关系类型"] == "上下级 / 盟友", detail=str(relation_payload)) + check("relation-repair-本窗演变稳定合并", + relation_payload["字段"]["演变轨迹"] == ["[窗8] 并肩作战;接受指挥"], + detail=str(relation_payload)) + check("relation-repair-其他字段并集完整", + relation_payload["字段"]["当前状态"] == "互信" + and relation_payload["字段"]["权力结构"] == "乙指挥甲", + detail=str(relation_payload)) + check("relation-repair-二次机械合并stderr可追踪", + "pair=(11, 12)" in repair_log and "合并数=1" in repair_log, + detail=repair_log) check("relation-repair-合并后整窗成功", db.window_status == "done", detail=str(db.window_errors)) check("relation-repair-成功不触发整窗重试", undo_calls == 0, detail=f"undo_calls={undo_calls}") +def test_repaired_relation_merge_order_and_reverse_are_stable(): + """二次机械合并对互补反向关系输出完全确定的 payload。""" + + forward_rows = [ + {"甲方": 20, "乙方": 21, "关系类型": "师徒", "本窗演变": "首次指点", + "其他字段": {"序列": "第二对"}}, + {"甲方": 12, "乙方": 11, "关系类型": "上下级", "本窗演变": "接受指挥", + "其他字段": {"权力结构": "乙指挥甲"}}, + {"甲方": 11, "乙方": 12, "关系类型": "盟友", "本窗演变": "并肩", + "其他字段": {"状态": "互信"}}, + ] + reverse_rows = [ + {"甲方": 12, "乙方": 11, "关系类型": "盟友", "本窗演变": "并肩", + "其他字段": {"状态": "互信"}}, + {"甲方": 21, "乙方": 20, "关系类型": "师徒", "本窗演变": "首次指点", + "其他字段": {"序列": "第二对"}}, + {"甲方": 11, "乙方": 12, "关系类型": "上下级", "本窗演变": "接受指挥", + "其他字段": {"权力结构": "乙指挥甲"}}, + ] + first = pu._merge_repaired_relation_output(forward_rows, {11, 12, 20, 21}) + second = pu._merge_repaired_relation_output(reverse_rows, {11, 12, 20, 21}) + expected = [ + { + "甲方": 11, + "乙方": 12, + "关系类型": "上下级 / 盟友", + "本窗演变": "并肩;接受指挥", + "其他字段": {"权力结构": "乙指挥甲", "状态": "互信"}, + }, + { + "甲方": 20, + "乙方": 21, + "关系类型": "师徒", + "本窗演变": "首次指点", + "其他字段": {"序列": "第二对"}, + }, + ] + check("relation-repair-互补反向输出完全相等", + first == second, detail=f"first={first}, second={second}") + check("relation-repair-确定pair片段和字段键顺序", + first == expected and list(first[0]["其他字段"]) == ["权力结构", "状态"], + detail=f"first={first}") + + +def test_repaired_relation_merge_rejects_unknown_top_level_dimensions(): + """repair 的非空未知顶层维度必须失败关闭;空值未知键可忽略。""" + + with_unknown = [{ + "甲方": 11, + "乙方": 12, + "关系类型": "盟友", + "本窗演变": "并肩", + "其他字段": {}, + "正文证据": "两人共同迎敌", + }] + try: + pu._merge_repaired_relation_output(with_unknown, {11, 12}) + except pu.RelationOutputConflict as exc: + check("relation-repair-非空未知顶层键失败关闭", + "正文证据" in str(exc), detail=str(exc)) + else: + check("relation-repair-非空未知顶层键失败关闭", False, + detail="未知顶层维度被静默丢弃") + + empty_unknowns = [{ + "甲方": 12, + "乙方": 11, + "关系类型": "盟友", + "本窗演变": "并肩", + "其他字段": {"状态": "互信"}, + "空字符串": " ", + "空对象": {}, + "空数组": [], + "空值": None, + }] + check("relation-repair-空值未知顶层键可忽略", + pu._merge_repaired_relation_output(empty_unknowns, {11, 12}) == [{ + "甲方": 11, "乙方": 12, "关系类型": "盟友", "本窗演变": "并肩", + "其他字段": {"状态": "互信"}, + }]) + + +def test_run_relation_conflict_repair_single_unknown_top_level_key_fails_window_without_relation_write(): + """首轮冲突后的单条 repair 也必须经过机械合并,未知键不得进入关系写事务。""" + + first = [ + {"甲方": 11, "乙方": 12, "关系类型": "盟友", "本窗演变": "继续合作", + "其他字段": {"当前状态": "互信"}}, + {"甲方": 12, "乙方": 11, "关系类型": "对手", "本窗演变": "公开决裂", + "其他字段": {"当前状态": "敌对"}}, + ] + repaired = [{ + "甲方": 11, + "乙方": 12, + "关系类型": "盟友", + "本窗演变": "继续合作", + "其他字段": {"当前状态": "互信"}, + "正文证据": "两人共同迎敌", + }] + db, prompts, update_calls, insert_calls, undo_calls, second_failure = _run_relation_repair_case( + first, repaired, return_error=True + ) + relation_cards = [ + row for row in db.drafts.values() + if row["payload"].get("type") == pu.RELATION_TYPE + ] + check("relation-repair-单条未知键CLI非零", second_failure is not None, + detail=str(second_failure)) + check("relation-repair-单条未知键每次只补一次重整", len(prompts) == 4, + detail=f"calls={len(prompts)}") + check("relation-repair-单条未知键计算后零关系写入", + update_calls == 0 and insert_calls == 0 and not relation_cards, + detail=f"update={update_calls}, insert={insert_calls}, cards={relation_cards}") + check("relation-repair-单条未知键最终整窗失败", + db.window_status == "failed" and undo_calls == 2, + detail=f"status={db.window_status}, undo_calls={undo_calls}") + + def test_run_relation_conflict_repair_conflict_fails_window_without_relation_write(): """关系专用重整仍冲突时不放宽比较,整窗重试后仍零关系写并失败关闭。""" @@ -3189,12 +3401,18 @@ class _LegacyRecoveryConn: self.embeddings = [{"draft_id": 701, "deleted": True, "entity_id": None}] if non_reset_deleted == "draft": self.drafts[0]["updater"] = "upgrade" + elif non_reset_deleted == "draft-undo": + self.drafts[0]["updater"] = "upgrade-undo" elif non_reset_deleted == "embedding-draft": self.drafts[0]["updater"] = "upgrade" + elif non_reset_deleted == "embedding-draft-undo": + self.drafts[0]["updater"] = "upgrade-undo" elif non_reset_deleted == "draft-status": self.drafts[0]["status"] = "confirmed" elif non_reset_deleted == "alias": self.aliases[0]["updater"] = "upgrade" + elif non_reset_deleted == "alias-undo": + self.aliases[0]["updater"] = "upgrade-undo" if active_artifact == "draft": self.drafts[0]["deleted"] = False elif active_artifact == "embedding": @@ -3243,11 +3461,12 @@ class _LegacyRecoveryConn: return _CardResult((sum(row["window_no"] == 7 for row in self.audits),)) if "example_upgrade_alias" in normalized: self._assert_artifact_sql( - "aliases", normalized, "a.deleted=TRUE", "a.updater='upgrade-reset'", + "aliases", normalized, "a.deleted=TRUE", "a.updater = ANY(%s)", ) + allowed = self._allowed_tombstone_updaters(params) return _CardResult((sum( row["work_id"] == 8 and row["window_no"] == 7 - and not (row["deleted"] and row["updater"] == "upgrade-reset") + and not (row["deleted"] and row["updater"] in allowed) for row in self.aliases ),)) if "example_upgrade_presence" in normalized: @@ -3261,13 +3480,14 @@ class _LegacyRecoveryConn: if "example_knowledge_embedding" in normalized: self._assert_artifact_sql( "embeddings", normalized, "e.entity_id IS NULL", "e.deleted=TRUE", - "d.deleted=TRUE", "d.updater='upgrade-reset'", "d.status='pending'", + "d.deleted=TRUE", "d.updater = ANY(%s)", "d.status='pending'", ) + allowed = self._allowed_tombstone_updaters(params) return _CardResult((sum( row["draft_id"] == 701 and not ( row["entity_id"] is None and row["deleted"] and self.drafts[0]["deleted"] - and self.drafts[0]["updater"] == "upgrade-reset" + and self.drafts[0]["updater"] in allowed and self.drafts[0]["status"] == "pending" ) for row in self.embeddings @@ -3285,13 +3505,14 @@ class _LegacyRecoveryConn: ),)) if "muse_knowledge_draft" in normalized: self._assert_artifact_sql( - "new_cards", normalized, "d.deleted=TRUE", "d.updater='upgrade-reset'", + "new_cards", normalized, "d.deleted=TRUE", "d.updater = ANY(%s)", "d.status='pending'", ) + allowed = self._allowed_tombstone_updaters(params) return _CardResult((sum( row["work_id"] == 8 and row["source_type"] == pu.SOURCE_TYPE and row["origin"] == "升格@窗7" and not ( - row["deleted"] and row["updater"] == "upgrade-reset" + row["deleted"] and row["updater"] in allowed and row["status"] == "pending" ) for row in self.drafts @@ -3313,6 +3534,86 @@ class _LegacyRecoveryConn: if missing: raise AssertionError(f"{name} 计数 SQL 缺少过滤条件:{missing}; sql={normalized}") + @staticmethod + def _allowed_tombstone_updaters(params): + """从生产 SQL 的数组参数读取本次上下文允许的墓碑 updater。""" + + for value in params: + if isinstance(value, list): + return set(value) + raise AssertionError(f"产物计数 SQL 未绑定 allowed_tombstone_updaters:params={params}") + + +class _CompensationConn(_LegacyRecoveryConn): + """无 marker 普通补偿夹具,复用七域计数并显式保留窗口状态。""" + + def __init__(self, non_reset_deleted=None): + super().__init__(non_reset_deleted=non_reset_deleted) + self.window_status = "pending" + self.window_error = None + + def execute(self, query, params=()): + normalized = " ".join(query.split()) + if normalized.startswith("SELECT id, status, error_message FROM example_upgrade_window"): + self.queries.append(normalized) + return _CardResult((91, self.window_status, self.window_error)) + result = super().execute(query, params) + if normalized.startswith("UPDATE example_upgrade_window SET status='failed'"): + self.window_status = "failed" + return result + + +def test_compensation_tombstone_domains_are_separate(): + """普通补偿只放行 alias/embedding/new_card 的 upgrade-undo 墓碑,其余产物仍阻断。""" + + for tombstone in ("draft-undo", "alias-undo", "embedding-draft-undo"): + clean = _CompensationConn(non_reset_deleted=tombstone) + with patch.object(pu.psycopg, "connect", return_value=clean): + pu._compensate_window(8, 7, "普通阶段失败") + check( + f"compensation-{tombstone}-upgrade-undo墓碑放行", + clean.window_status == "failed" + and str(clean.window_error).startswith(pu.RETRYABLE_CLEAN_PREFIX) + and pu.COMPENSATION_FAILED_PREFIX not in str(clean.window_error), + detail=str(clean.window_error), + ) + + for artifact, row in ( + ("presence", {"work_id": 8, "window_no": 7, "updater": "upgrade-undo"}), + ("card_state", {"draft_id": 701, "watermark_window": 7, "updater": "upgrade-undo"}), + ("audit", {"draft_id": 701, "window_no": 7, "updater": "upgrade-undo"}), + ): + blocked = _CompensationConn() + getattr(blocked, { + "presence": "presence", + "card_state": "card_states", + "audit": "audits", + }[artifact]).append(row) + with patch.object(pu.psycopg, "connect", return_value=blocked): + try: + pu._compensate_window(8, 7, "普通阶段失败") + except pu.CompensationFenceConflict: + pass + else: + check(f"compensation-{artifact}-看似undo仍阻断", False, + detail="未抛 CompensationFenceConflict") + check( + f"compensation-{artifact}-看似undo持久compensation-failed", + blocked.window_status == "failed" + and str(blocked.window_error).startswith(pu.COMPENSATION_FAILED_PREFIX), + detail=str(blocked.window_error), + ) + + invalid = _CompensationConn() + try: + pu._window_artifact_counts(invalid, 8, 7, ("upgrade",)) + except ValueError as exc: + check("compensation-非allowlist-updater拒绝", "不允许" in str(exc), detail=str(exc)) + else: + check("compensation-非allowlist-updater拒绝", False, + detail="非 allowlist updater 被接受") + check("compensation-非法updater拒绝不查库", invalid.queries == [], detail=str(invalid.queries)) + def test_recover_legacy_failed_cli_contract(): """legacy failed 只允许清理副作用为零的窗口,只有安全 reset 墓碑可豁免。""" @@ -3433,6 +3734,21 @@ def test_recover_legacy_failed_cli_contract(): detail=non_reset_preview.output, ) + for artifact in ("draft-undo", "embedding-draft-undo", "alias-undo"): + undo_conn = _LegacyRecoveryConn(non_reset_deleted=artifact) + with patch.object(pu.psycopg, "connect", return_value=undo_conn): + undo_preview = CliRunner().invoke( + pu.cli, ["recover-legacy-failed", "--work-id", "8", "--window-no", "7", "--preview"] + ) + undo_data = json.loads(undo_preview.output) + key = {"draft-undo": "new_cards", "embedding-draft-undo": "embeddings", + "alias-undo": "aliases"}[artifact] + check( + f"recover-legacy-failed-{artifact}-默认仍阻断", + undo_preview.exit_code == 0 and undo_data["artifact_counts"][key] > 0, + detail=undo_preview.output, + ) + deleted_presence_conn = _LegacyRecoveryConn(active_artifact="presence-deleted") with patch.object(pu.psycopg, "connect", return_value=deleted_presence_conn): deleted_presence_preview = CliRunner().invoke( @@ -3473,11 +3789,15 @@ if __name__ == "__main__": test_merge_card_chapter_evidence, test_run_alias_paths_store_real_canonical_name, test_run_new_card_registers_normalized_name_and_aliases_same_window, + test_new_card_same_window_update_uses_database_initial_revision, test_run_exact_and_disjoint_duplicate_entity_updates_merge_once, test_run_conflicting_duplicate_entity_updates_fail_before_any_write, test_run_relation_cards_enter_touched_embedding, test_run_conflicting_duplicate_relations_fail_before_any_relation_write, test_run_relation_conflict_repair_merges_to_one_card, + test_repaired_relation_merge_order_and_reverse_are_stable, + test_repaired_relation_merge_rejects_unknown_top_level_dimensions, + test_run_relation_conflict_repair_single_unknown_top_level_key_fails_window_without_relation_write, test_run_relation_conflict_repair_conflict_fails_window_without_relation_write, test_run_identical_duplicate_relations_write_existing_card_once, test_run_prompt_revision_rejects_entity_and_relation_stale_results, @@ -3498,6 +3818,7 @@ if __name__ == "__main__": test_capture_window_input_rejects_incomplete_or_empty_content, test_two_phase_external_calls_and_atomic_embedding_contract, test_recovery_and_redo_rejection_contracts, test_real_pg_rollback_smoke_entry, + test_compensation_tombstone_domains_are_separate, test_recover_legacy_failed_cli_contract, test_capture_window_state_binds_current_window_and_preserves_other_status, test_load_window_material_orders_blocks_like_capture_input, diff --git a/docs/2026-07-16-升格卡改造设计.md b/docs/2026-07-16-升格卡改造设计.md index f7d7701..cb359d0 100644 --- a/docs/2026-07-16-升格卡改造设计.md +++ b/docs/2026-07-16-升格卡改造设计.md @@ -249,8 +249,8 @@ flowchart TB 升格抽取不得在模型或嵌入调用期间持有业务数据库连接。一次窗口只保留一个可恢复的中间提交点,避免把每个更新批次都变成独立恢复协议: 1. **无连接计算**:读取正文、既有卡和 revision 后立即关闭连接;完成观察、证据修复、语义判重和全部实体更新输出。此阶段失败时没有知识写入,只把窗口留成干净失败断点。 -2. **实体短写事务**:先复验模型输入快照,再复验模型所见 revision,原子写入新卡、别名、留档和实体更新;把窗口置为 `processing`,并在同一事务记录本书受控域的完整状态摘要。进程此后崩溃,续跑必须先复验摘要,再精确撤销本窗。 -3. **无连接关系计算**:实体阶段提交后重新读取人物与关系快照,关闭连接,再生成或重整关系。卡片已更新后的状态仍能进入关系判断,但模型调用期间没有长事务。 +2. **实体短写事务**:先复验模型输入快照,再复验模型所见 revision,原子写入新卡、别名、留档和实体更新;新卡 `INSERT` 同时返回数据库实际 `revision`,同窗后续更新的 CAS 必须绑定该值,不能假定初始版本。把窗口置为 `processing`,并在同一事务记录本书受控域的完整状态摘要。进程此后崩溃,续跑必须先复验摘要,再精确撤销本窗。 +3. **无连接关系计算**:实体阶段提交后重新读取人物与关系快照,关闭连接,再生成或重整关系。首次同 pair 冲突只调用一次模型 repair;repair 仍按同 pair 拆条时,只允许机械合并关系类型、本窗演变和无字段冲突的其他维度,非空未知顶层键或同字段异值继续失败关闭,不调用第三次模型。卡片已更新后的状态仍能进入关系判断,但模型调用期间没有长事务。 4. **关系短写事务**:复验 `processing` 摘要、输入快照、人物 revision 和关系集合,写关系与剩余出场章;窗口仍保持 `processing`,原子刷新状态摘要。 5. **严格嵌入与完成事务**:在无业务数据库连接状态下计算本窗 touched 卡向量,再用短事务复验摘要和向量 owner 后落库;只有嵌入成功且摘要再次一致,窗口才置 `done`。普通嵌入失败也按窗口失败处理,不允许后窗在缺向量状态下继续语义判重。 @@ -273,8 +273,9 @@ flowchart LR - `inputSha256` 绑定实际送模输入:作品标题、窗口边界、每章 `id/title/status/revision`、每个 block 的 `id/order/type/revision/content_text`,以及实体 schema active version/字段合同、抽取合同和当前 `parse_upgrade.py` SHA-256。按章拼接正文的规则必须与送模文本一致;实体阶段写入前与关系阶段写入前都要重算,章域不完整、无 block 或拼接正文为空直接失败关闭。 - `stateSha256` 覆盖本书卡、别名、留档、水位、审计、向量和全部窗口不变元数据;窗口自身可变的状态/错误文本不参与自引用 hash。`processing` 恢复和失败补偿都要在固定顺序的七域锁内复验它。 - 任何外部确认、改版、删除或辅助域漂移都写 `compensation-failed` 并停止,禁止覆盖外部状态。嵌入落库与 `done` 在同一最终短事务完成,消除 `done` 已提交但向量尚未落库的崩溃缝隙。 +- 普通无 marker clean 的产物检查只豁免完整满足墓碑条件且 `updater` 为 `upgrade-reset` 或 `upgrade-undo` 的 alias/new_card/embedding;presence、card_state、audit 仍有行即阻断。该上下文用于识别 `undo_window` 合法留下的新卡墓碑,不得放宽 status、entity owner 或其他 updater。 - 当前数据库没有所有写方共同遵守的 work 级互斥键。为保证摘要无幻读,七个受控域使用**只跨短写事务**的固定表锁;它会短暂串行化不同作品的写阶段,但绝不跨模型或嵌入计算。待真实耗时证明成为瓶颈后,再单独设计全系统统一的 work 级写锁,不能在本次修复里用弱锁换假安全。 -- fence 引入前的 `legacy failed` 窗不能自动猜测。`recover-legacy-failed` 是唯一恢复入口:preview 输出窗口、七域本窗产物计数与确认 hash;execute 必须取得同书锁、精确匹配 hash、确认无 `processing` 窗且**阻断产物计数为 0**,才可标成 `retryable-clean`。计数只豁免能机械证明来源的全书 reset 墓碑:draft 必须同时 `deleted=true`、`updater='upgrade-reset'`、`status='pending'`;embedding 还必须 `entity_id IS NULL` 且自身 `deleted=true`;alias 必须 `deleted=true` 且 `updater='upgrade-reset'`。presence 没有 reset 来源标记,任何本窗行(包括软删墓碑)都阻断;其他软删、活跃行、已有 entity owner,以及任何 card_state/audit 行也都阻断。完整 `stateSha` 仍覆盖所有 deleted 行,不能混淆。操作人仍须先确认旧进程与数据库 backend 已结束,工具不得把“查不到已提交行”误当成“旧事务不存在”。 +- fence 引入前的 `legacy failed` 窗不能自动猜测。`recover-legacy-failed` 是唯一恢复入口:preview 输出窗口、七域本窗产物计数与确认 hash;execute 必须取得同书锁、精确匹配 hash、确认无 `processing` 窗且**阻断产物计数为 0**,才可标成 `retryable-clean`。计数只豁免能机械证明来源的全书 reset 墓碑:draft 必须同时 `deleted=true`、`updater='upgrade-reset'`、`status='pending'`;embedding 还必须 `entity_id IS NULL` 且自身 `deleted=true`;alias 必须 `deleted=true` 且 `updater='upgrade-reset'`。`upgrade-undo` 只属于普通 clean,在 legacy 上仍必须阻断,两种上下文不得混用。presence 没有 reset 来源标记,任何本窗行(包括软删墓碑)都阻断;其他软删、活跃行、已有 entity owner,以及任何 card_state/audit 行也都阻断。完整 `stateSha` 仍覆盖所有 deleted 行,不能混淆。操作人仍须先确认旧进程与数据库 backend 已结束,工具不得把“查不到已提交行”误当成“旧事务不存在”。 - 显式 `--redo-window` 需要把 redo 前完整快照持久化到进程外,单靠内存快照无法承受崩溃。该持久快照协议另立设计前,命令机械拒绝;历史修正走已有全书 backup/reset/rebuild,不提供不完整的 redo。 ### 风险