修复: 收紧升格关系重整与补偿边界

This commit is contained in:
zizi 2026-07-23 11:49:26 +08:00
parent c54d188887
commit 0446a9e8f0
4 changed files with 511 additions and 41 deletions

View File

@ -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。

View File

@ -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'

View File

@ -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,

View File

@ -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。
### 风险