zizi 46704ee5da fix(parse-book): 独立复核揪出的2个真bug——recluster跨书实例归属被覆写+迁移章号正则误抓数量词
Bug1(parse_ingest _absorb):归并时用输家书名无条件覆写每条实例的来源书名→'二级归并母卡'里跨书实例的真实来源被抹成输家书名;改为保留实例原有书名(ins.get('书') or book)。已跑的718归并中75条实例/31母卡受影响,可从回溯归并审计的被并卡全文还原(待remediate)。

Bug2(parse_upgrade INLINE_CHAP_RE):裸'N章'分支把'隔3章/花了3章篇幅'数量词误当章号、还标最高置信'内嵌',违背人工精确化(拍板#5);改为只认'第X章'/'Ch.X'(区间取首章'第489-490章'→489),裸写法一律落待LLM重抽。迁移仅dry-run未落库、无数据影响。11条边界用例亲测通过。

独立opus子代理对抗复核确认无阻断性问题:无物理删知识卡/无dry-run漏写/无额度绕过/软删可逆/连接三段/调用契约完好。
2026-07-17 11:52:54 +08:00

1142 lines
69 KiB
Python
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 python3
"""parse-book skill:作品面升格执行器(方案 v6,创始人 2026-07-14 认可)。
把参考书自己的实体/关系/大纲按「全实体统一生长」机制抽进正式候选层:
一实体一卡(不分全卡/轻卡),每窗流程=
⓪ 机械预扫(零 AI):全库已知名字在本窗正文精确扫描 → 在场已知实体子集
① 实体观察(AI×1):新名字(顺带产初卡)/ 已知实体新信息 / 纯出场登记
② 判重(机械+观察自判):别名表/留档表精确查 → 疑似别名转观察材料
③ 立卡门槛:跨章戏份才立卡(单章龙套进出场留档,跨窗合计跨章再补立)
④ 卡更新(AI×0-3):只对有新信息的卡,每批≤6 张,读全卡→只回变更字段→三类合并
⑤ 关系增量(AI×1):核心角色两两关系变化 → 滚进关系卡(甲乙锚草稿编号)
⑥ 机械收尾:出场留档/窗状态/卡水位/覆写审计
防膨胀(第四轮评审 G1/G5):观察调用只带「窗内命中」实体索引(不带全库);
窗切割限 12 章/3.5 万字双闸。防重跑自噬:同窗重跑先按审计撤销再重写。
嵌入延后:试跑期不调嵌入 API(判重近邻延后批量做),不混烧请求额度。
命令:
windows --work-id N 机械切正文窗(幂等,from_chapter 锚)
run --work-id N [--max-windows K] [--max-calls M] 按窗顺序跑,断点续跑
status --work-id N 进度
"""
import json
import pathlib
import re
import sys
import click
import psycopg
# 复用章级管线的敏感降级链与 llm 入口(trust_env/重试/JSON 容错同源)
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parent))
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "llm" / "scripts"))
from parse_llm import m3_json, SensitiveHardStop, IDENTITY, TENANT, DSN # noqa: E402
# 语义判重(P1)复用 embed skill 的嵌入通道(同模型同维、与检索端语义对齐)——
# 只在开启 --semantic-dedup 时才真调,默认关(试跑期嵌入延后,见文件头注释)
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "embed" / "scripts"))
from embed_drafts import _session as _embed_session, embed_texts # noqa: E402
# ── 窗切割参数(方案 §五B:3–5 万字/窗、10–15 章,取保守双闸防观察输出过载)──
WIN_MAX_CHARS = 35000 # 单窗正文字数上限
WIN_MAX_CHAPS = 12 # 单窗章数上限
UPDATE_BATCH = 6 # 卡更新每批张数上限(防长清单丢字段)
SOURCE_TYPE = "upgrade_book" # 独立来源标记:与范式卡 parse_book 隔离判重/检索/确认
# 作品面七型(六实体型+关系型;大纲卡全书收尾单独做,不在窗循环内)
ENTITY_TYPES = ("character", "location", "item", "faction", "power_system", "event")
RELATION_TYPE = "character_relation"
# 追加类字段白名单(值为数组的字段一律追加)。「演变历程」是里程碑对象数组,走独立对象
# 合并路径(MILESTONE_FIELDS);其余是字符串条目数组(历史上带 [窗N] 前缀)。「大事记」「经历」
# 是历史孤儿名(任何 schema 都没定义、会被合同守卫裁掉),保留仅为向后兼容旧数据、不再新用——
# 升格卡改造(2026-07-17)后已发生台阶统一记入「演变历程」。
APPEND_FIELDS = {"成长弧线", "演变轨迹", "大事记", "经历", "演变历程"}
# 里程碑对象数组字段:条目是 {章,台阶,周期} 结构化对象,去重按台阶内容、排序按真实章号
# (不用运行时窗号)——升格卡改造 P0 落点,区别于上面的字符串条目追加字段。
MILESTONE_FIELDS = {"演变历程"}
# 生命周期枚举(设计稿 §4.2):每条里程碑「周期」的取值域。
LIFECYCLE = ("登场", "成长", "高光", "退场", "结局")
# 语义判重(P1,设计稿 §8.2):召回同书近邻相似度 ≥ 此阈值才交 M3 终判。
DEDUP_SIM_THRESHOLD = 0.78
CHAP_BIG = 10 ** 9 # 章号缺失/待人工的里程碑,排序时排到最后
# ── 库内合同(元数据驱动公理:prompt 与守卫同源,禁手写合同)──
def load_entity_contracts(conn):
"""加载作品面七型的字段合同(走 active_version_id,与章级管线同语义)。"""
contracts = {}
for t in ENTITY_TYPES + (RELATION_TYPE,):
snap = conn.execute(
"""SELECT v.field_contract_snapshot FROM muse_meta_schema_version v
JOIN muse_meta_schema s ON s.active_version_id=v.id
WHERE s.tenant_id=%s AND s.schema_key=%s""", (TENANT, t)).fetchone()[0]
contracts[t] = {"中文名": snap.get("中文名", t), "判据": snap.get("判据", ""),
"字段": [f for f in snap.get("特有字段", [])]}
return contracts
def render_entity_contracts(contracts, types):
"""合同渲染为 markdown 表(观察/更新 prompt 共用)。"""
parts = []
for t in types:
c = contracts[t]
rows = "\n".join(f"| {f['key']} | {f.get('说明', '')} |" for f in c["字段"])
parts.append(f"### {t}({c['中文名']})\n判据:{c['判据']}\n\n"
f"| 字段 key | 说明 |\n|---|---|\n{rows}")
return "\n\n".join(parts)
# ── 窗切割(机械,零 AI)──
def cut_windows(conn, work_id):
"""按章边界贪心切正文窗:累计超 3.5 万字或 12 章即断窗。幂等(from_chapter 锚)。"""
rows = conn.execute(
"""SELECT c.order_no, COALESCE(b.word_count, length(b.content_text))
FROM muse_content_chapter c
JOIN muse_content_block b ON b.chapter_id=c.id AND b.deleted=FALSE
WHERE c.tenant_id=%s AND c.work_id=%s AND c.deleted=FALSE
ORDER BY c.order_no""", (TENANT, work_id)).fetchall()
wins, cur, chars = [], [], 0
for order_no, wc in rows:
if cur and (chars + (wc or 0) > WIN_MAX_CHARS or len(cur) >= WIN_MAX_CHAPS):
wins.append((cur[0], cur[-1]))
cur, chars = [], 0
cur.append(order_no)
chars += (wc or 0)
if cur:
wins.append((cur[0], cur[-1]))
n = 0
for i, (a, b) in enumerate(wins, 1):
r = conn.execute(
"""INSERT INTO example_upgrade_window
(work_id, window_no, from_chapter, to_chapter, tenant_id)
VALUES (%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, work_id, from_chapter) DO NOTHING""",
(work_id, i, a, b, TENANT))
n += r.rowcount
conn.commit()
return len(wins), n
# ── 每窗材料与已知名加载 ──
def load_window_text(conn, work_id, a, b):
"""窗内正文拼接:第N章《标题》+正文。"""
rows = conn.execute(
"""SELECT c.order_no, c.title, b2.content_text
FROM muse_content_chapter c
JOIN muse_content_block b2 ON b2.chapter_id=c.id AND b2.deleted=FALSE
WHERE c.tenant_id=%s AND c.work_id=%s AND c.order_no BETWEEN %s AND %s
AND c.deleted=FALSE ORDER BY c.order_no""",
(TENANT, work_id, a, b)).fetchall()
return "\n\n".join(f"## 第{o}章 {t}\n{x}" for o, t, x in rows)
def load_known(conn, work_id):
"""加载判重底册:name_map(名字/别名→卡)+ presence 留档出场章。
返回 (name_map: {名字: (draft_id, 型, 摘要)}, presence: {(型,名字): set(章)})"""
name_map = {}
for did, payload in conn.execute(
"""SELECT id, draft_payload FROM muse_knowledge_draft
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE""",
(TENANT, work_id, SOURCE_TYPE)).fetchall():
t, nm = payload.get("type", ""), (payload.get("名称") or "").strip()
brief = payload.get("一句话摘要", "")
if nm:
name_map[nm] = (did, t, brief)
for al in payload.get("别名", []) or []:
if al and al.strip():
name_map[al.strip()] = (did, t, brief)
for cn, al in conn.execute(
"SELECT canonical_name, alias FROM example_upgrade_alias "
"WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE",
(TENANT, work_id)).fetchall():
if cn in name_map:
name_map[al] = name_map[cn]
presence = {}
for t, nm, ch in conn.execute(
"SELECT entity_type, name, chapter_no FROM example_upgrade_presence "
"WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE",
(TENANT, work_id)).fetchall():
presence.setdefault((t, nm), set()).add(ch)
return name_map, presence
def prescan(name_map, text):
"""机械预扫(G1 防膨胀核心):只把「本窗正文出现」的已知实体带进观察调用。"""
hit = {}
for nm, (did, t, brief) in name_map.items():
if nm and nm in text:
# 同卡多名字只留一条(正名优先:先插入的是正名)
hit.setdefault(did, (nm, t, brief))
return {nm: (did, t, brief) for did, (nm, t, brief) in hit.items()}
# ── 提示词(缓存友好:固定规则前置,正文窗次之且窗内多调用共享,任务尾置)──
def observe_prompt(contracts, title, a, b, text, onstage):
onstage_lines = "\n".join(f"- {t}|{nm}|{brief}" for nm, (_, t, brief) in onstage.items()) or "(无)"
return f"""【功能指令(parse-book 作品面升格·实体观察)】
通读本窗正文,产出三类结果(只输出一个 JSON 对象):
1) 新名字:正文出现、但「在场已知实体」清单里没有的实体(六型:{"/".join(ENTITY_TYPES)})。每个给:型、名称、别名、一句话摘要、按该型合同能填的字段(有正文证据才填)、出场章号。若你怀疑它其实是清单中某已知实体的别名/改名/化名,填「疑似别名指向」。若本窗已见其登场或首个进化台阶,按合同给「演变历程」的首条里程碑对象。
2) 已知实体新信息:清单中实体在本窗的实质新信息(境界变化/性格显露/重大经历/立场转变),一条 ≤60 字观察点。没有实质新信息的不要报。
3) 纯出场:清单中实体本窗出现但无实质新信息的,只报名称+出场章。
纪律:一次性龙套(单章无名或仅路过)不报进新名字;实体判据与字段以合同为准,无证据不填;不脑补。数值(战力/指数/排名等)必须正文原样出现才可写,禁止推算或编造。组织改组/合并产生的新组织是**新实体**(走新名字),不是旧组织的别名。「疑似别名指向」只在确为同一实体改名/化名时填。
里程碑纪律(凡填「演变历程」字段必守):每条是对象 {{"章": 正文原样出现的真实章号(整数如 420,跨多章连续事件用区间字符串如 "420-423"), "台阶": "进化到什么+靠什么事件的一句话", "周期": 登场/成长/高光/退场/结局 之一}};**必须用真实章号,禁止 [窗N] 窗号、禁止"本窗/近期/前段"这类相对指代**(卡会脱离运行环境被单独阅读);抽不出完整对象时至少给 {{"章","台阶"}},别整条丢;当前态字段(品阶/能力与限制/摘要等)只写"现在什么样"的干净值,历史进化流水一律进「演变历程」、不许塞进当前态字段。
【六型字段合同】
{render_entity_contracts(contracts, ENTITY_TYPES)}
【输出规则(只输出一个 JSON 对象)】
{{"新名字": [{{"型": "character", "名称": "", "别名": [], "一句话摘要": "", "字段": {{}}, "出场章": [章号], "疑似别名指向": ""}}],
"已知实体新信息": [{{"名称": "", "观察点": "", "出场章": [章号]}}],
"纯出场": [{{"名称": "", "出场章": [章号]}}]}}
━━━ 本窗材料(每窗不同,非规则)━━━
《{title}》第 {a}–{b} 章正文:
{text}
【在场已知实体】(预扫命中,判重参照)
{onstage_lines}"""
def update_prompt(contracts, title, a, b, text, cards_with_obs):
cards_json = json.dumps([{"draft_id": d, "当前卡": p, "本窗观察点": o}
for d, p, o in cards_with_obs], ensure_ascii=False, indent=1)
return f"""【功能指令(parse-book 作品面升格·卡增量更新)】
下列每张卡给出「当前卡全文」与「本窗观察点」。对照本窗正文,**只输出需要变更的字段**:
- 覆写类字段(性格底色/说话方式/当前状态等标量):**必须输出该字段完整的新全量值**——旧值里仍然成立的信息要保留进新值,禁止只写"新增…"式增量(那会把旧信息抹掉);
- 里程碑字段(演变历程):每当实体发生境界/代际/形态/能力的跃迁,或到达登场/高光/退场/结局节点,**追加**一条里程碑对象 {{"章": 真实章号(整数如 420 或跨章区间字符串 "420-423"), "台阶": "进化到什么+靠什么事件的一句话", "周期": 登场/成长/高光/退场/结局 之一}}——只输出本窗**新增**里程碑(不重抄旧条目,不输出 _win 等内部键);章必须是正文原样章号、**禁 [窗N] 与"本窗/近期"相对指代**;抽不出完整对象时至少给 {{"章","台阶"}},别整条丢;
- 其他追加类字段(成长弧线=未来计划、演变轨迹等数组):只输出本窗新增条目(不要重抄旧条目,不要自己加 [窗N] 前缀,系统会加);
- **当前态字段干净纪律**:品阶/能力与限制/当前状态/一句话摘要等只写"现在是什么样"的全量当前值——历史进化流水一律进「演变历程」里程碑,**绝不把成长史塞进当前态字段**(这是老卡把升级线塞进「戏剧作用/流转计划/跨体系换算」污染字段语义的病根,务必避免);
- 没有变化的字段不要输出;整卡无实质变化则不输出该卡;
- 该实体的身份/处境已发生重大变化时,把「一句话摘要」也作为变更字段输出(写当前态全量)。
纪律:以正文为证据,不脑补;数值必须正文原样出现,禁止推算;字段 key 必须来自该型合同;**禁止「本窗/本段/近期」等相对时间指代**——写绝对状态或带章号(卡会脱离本窗被单独阅读);只写**该实体自己**的信息(他人的任命/心理/变化不得写入本卡);每卡别名有新发现可在「别名新增」里给——只收**该实体自己**的新别名(他人对它的称呼算,它对别的实体的称呼不算),且必须是可在正文原样出现的纯名字(禁带括号注释与说明文字,禁单字)。
【相关型字段合同】
{render_entity_contracts(contracts, sorted({p.get("type") for _, p, _ in cards_with_obs} & set(ENTITY_TYPES)))}
【输出规则(只输出一个 JSON 对象;即使本批所有卡都无变化,也必须输出 {{"更新": []}},不得省略"更新"键)】
{{"更新": [{{"draft_id": 数字, "变更字段": {{"字段key": "新值或新增条目数组"}}, "别名新增": []}}]}}
━━━ 本窗材料 ━━━
《{title}》第 {a}–{b} 章正文:
{text}
【待更新的卡】
{cards_json}"""
def relation_prompt(contracts, title, a, b, text, char_cards, existing_rels):
chars = "\n".join(f"- draft_id={d}|{p.get('名称')}|{p.get('一句话摘要', '')}"
for d, p in char_cards)
rels = "\n".join(f"- {r.get('甲方名称')} × {r.get('乙方名称')}:{r.get('关系类型', '')}"
for _, r in existing_rels) or "(暂无)"
return f"""【功能指令(parse-book 作品面升格·人物关系增量)】
基于本窗正文,报告下列核心角色**两两之间**的关系变化(新建立的关系 / 已有关系的演变)。
只报有正文证据的实质变化;没有变化输出空数组。
【character_relation 字段合同】
{render_entity_contracts(contracts, [RELATION_TYPE])}
【输出规则(只输出一个 JSON 对象;甲乙用 draft_id 指认)】
{{"关系": [{{"甲方": 数字, "乙方": 数字, "关系类型": "", "本窗演变": "≤60字", "其他字段": {{}}}}]}}
━━━ 本窗材料 ━━━
《{title}》第 {a}–{b} 章正文:
{text}
【核心角色】
{chars}
【已有关系(避免重报建立)】
{rels}"""
# ── 合并与审计(三类字段演进 + 撤销依据)──
WIN_PREFIX_RE = None # 延迟编译(模块顶部 import re 已有)
def _strip_prefix(s):
"""剥条目行首的窗号类前缀(可能多重堆叠;含模型自造的 [窗本窗] 等变体),返回纯内容。"""
import re as _re
global WIN_PREFIX_RE
if WIN_PREFIX_RE is None:
# [窗…]/[本窗…] 任意变体全剥(窗29实测模型自造 "[窗本窗]",仅数字版剥不掉;
# 深空实测又造 "[窗387-388]" 章号范围变体,7 字符超旧上限 6——放宽到 12)
WIN_PREFIX_RE = _re.compile(r"^(?:\[[窗本][^\]]{0,12}\]\s*)+")
return WIN_PREFIX_RE.sub("", str(s)).strip()
TAIL_DEBRIS_RE = None # 尾部 JSON 拼接残渣(延迟编译)
def _strip_tail(s):
"""剥条目尾部的 JSON 拼接残渣(批9 样张走查实证:李锋等 20 张卡 59 处条目
尾挂 ", / '] / "} 等符号——模型把结构化输出的收尾符号带进了条目文本)。反复剥直到干净。"""
import re as _re
global TAIL_DEBRIS_RE
if TAIL_DEBRIS_RE is None:
TAIL_DEBRIS_RE = _re.compile(r"""(?:",|'\]|"\]|"}|'}|',)\s*$""")
t = str(s).rstrip()
while True:
m = TAIL_DEBRIS_RE.search(t)
if not m:
return t
t = t[:m.start()].rstrip()
def _clean_alias(al):
"""别名机械准入(抽检#6):拒收含括号注释/超长的备忘录式别名——预扫精确匹配永不命中=死数据。
单字别名一并拒收(复检 M5:'新'字入表后预扫全文命中率爆炸,纯噪声)。"""
al = (al or "").strip()
if not al or len(al) < 2 or len(al) > 12:
return None
if any(c in al for c in "()()。,,:"):
return None
return al
def _is_garbage(text):
"""结构垃圾检测(抽检 H2 根治):窗75 实测模型把原始变更包字符串化塞进字段值,
7 张卡被 {'draft_id':...} 类程序结构污染。含结构特征的文本一律拒收留审计。"""
t = str(text)
return ("draft_id" in t or "变更字段" in t or "别名新增" in t
or t.lstrip().startswith(("{'", '{"', "[{")))
def _entry_win(x):
"""追加条目的窗号(无 [窗N] 前缀的初卡条目记 0,排最前)。"""
m = re.match(r"^\[窗(\d+)\]", str(x))
return int(m.group(1)) if m else 0
# ── 里程碑对象(升格卡改造 P0):{章,台阶,周期} 结构化条目的清洗/合并 ──
# 章级排序不用运行时窗号,而用真实章号(设计稿 §6.1「真实章号索引」的落点)。
# 内嵌真实章号抽取:正文/台阶里原样出现的「第X章 / X章 / Ch.X」——迁移与降级共用
# 只认"第X章"/"Ch.X"两种明确章号写法。刻意不收裸"N章"——它无法与"隔3章/花了3章篇幅"这类
# 数量词区分,误当章号会污染迁移(且被标成最高置信"内嵌"),违背人工精确化;裸写法一律落到待LLM重抽。
# 带"第"前缀时取首个数字,天然处理"第489-490章"→489(区间取首章、不取末章)。
INLINE_CHAP_RE = re.compile(r"第\s*(\d+)(?:\s*[-—~到]\s*\d+)?\s*章|(?:Ch|CH|ch)\.?\s*(\d+)")
def _extract_inline_chapter(text):
"""从台阶文本抽第一个内嵌真实章号(第X章/ChX/X章),抽不到返回 None。"""
m = INLINE_CHAP_RE.search(str(text))
if not m:
return None
return int(next(g for g in m.groups() if g))
def _chapter_sort_key(ch):
"""里程碑排序键:从「章」值(整数 / 区间字符串"420-423" / None)取起点章号;
缺失或非法排到最后(CHAP_BIG),保证"待人工"条目不插进正常时间线中间。"""
if isinstance(ch, int):
return ch
if isinstance(ch, str):
m = re.search(r"\d+", ch)
if m:
return int(m.group())
return CHAP_BIG
def _infer_lifecycle(text):
"""从台阶文本启发式推断生命周期枚举(模型未给或非法「周期」时的降级填充)。
诚实边界:这是关键词启发式、非精确判定;新抽取由提示词强制模型直接给枚举,此路仅兜底。"""
t = str(text)
if any(w in t for w in ("登场", "首次", "初次", "出场", "诞生", "创立", "问世", "面世")):
return "登场"
if any(w in t for w in ("退役", "封存", "陨落", "覆灭", "销毁", "谢幕", "退场", "解散", "湮灭")):
return "退场"
if any(w in t for w in ("结局", "终局", "最终", "决战", "了结", "落幕")):
return "结局"
if any(w in t for w in ("巅峰", "高光", "对决", "突破", "解放", "觉醒", "封神", "登顶", "碾压")):
return "高光"
return "成长"
def _clean_milestone(item, win_no):
"""规范化单个里程碑为 {章,台阶,周期,_win};垃圾/空台阶返回 None。
- dict 入参:取 章/台阶/周期;台阶剥前缀+剥尾残;缺章从台阶抽内嵌章号兜底;缺/非法周期启发式推断。
- 字符串入参(模型降级输出或存量迁移):整串当台阶,抽内嵌章号当章,推断周期。
- _win 盖当前窗号,仅作撤销溯源(undo_window 按它删本窗新增),不参与展示与排序。
降级保底(设计稿 §九 风险2 + 拍板#1):抽不出完整对象也退成"章号+一句话"最小对象,绝不整条丢。"""
if isinstance(item, dict):
step = _strip_tail(_strip_prefix(str(item.get("台阶") or item.get("阶") or "")))
ch = item.get("章")
cycle = item.get("周期")
else:
step = _strip_tail(_strip_prefix(str(item)))
ch, cycle = None, None
if not step or _is_garbage(step):
return None
# 章号:对象已给(整数或含数字的区间字符串)就用;否则从台阶文本抽内嵌章号;再无则 None(待人工)
if not (isinstance(ch, int) or (isinstance(ch, str) and re.search(r"\d", ch))):
ch = _extract_inline_chapter(step)
if cycle not in LIFECYCLE:
cycle = _infer_lifecycle(step)
return {"章": ch, "台阶": step, "周期": cycle, "_win": win_no}
def _merge_milestones(old, items, win_no):
"""合并里程碑数组:去重按台阶内容、排序按真实章号。返回 (merged, rejected)。
old 中已有对象保留其原 _win(不被本窗覆盖);新对象由 _clean_milestone 盖当前 win_no。
rejected 为被判垃圾的原始条目,交调用方留审计(对齐字符串路径的垃圾拦截)。"""
kept = [m for m in (old or []) if isinstance(m, dict) and m.get("台阶")]
seen = {str(m.get("台阶", "")).strip() for m in kept}
rejected = []
for it in (items if isinstance(items, list) else [items]):
m = _clean_milestone(it, win_no)
if not m:
rejected.append(it)
continue
key = m["台阶"].strip()
if key and key not in seen:
kept.append(m)
seen.add(key)
return sorted(kept, key=lambda m: _chapter_sort_key(m.get("章"))), rejected
def merge_card(conn, draft_id, win_no, changes, alias_new, valid_keys=None):
"""按 5.1 三类规则合并变更字段:数组/白名单=追加(剥模型自带前缀+同文去重后带窗号),
标量=覆写留审计(增量式假全量拦截转追加——抽检#4 信息回退病)。
valid_keys:该型合同的合法字段 key 集——越合同 key 裁剪留审计(窗29实测模型把
整段条目文本误当字段 key 写入,无校验会把卡体字段区打烂)。"""
payload = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(draft_id,)).fetchone()[0]
fields = payload.setdefault("字段", {})
# 卡水位(该卡最后一次被更新的窗号):补跑迟到窗(如窗41在窗80后补跑)的覆写类字段
# 若直接落卡,会把书末态倒写回中期态(时间倒流污染,2026-07-15 补跑实测坐实)。
# 闸门:win_no < 水位 ⇒ 覆写只留审计不动卡体;追加类带窗号标签乱序无害,不拦。
wm_row = conn.execute("SELECT watermark_window FROM example_upgrade_card_state WHERE draft_id=%s",
(draft_id,)).fetchone()
watermark = wm_row[0] if wm_row else 0
changes = dict(changes or {})
# 模型偶把「别名新增」混进变更字段(窗2实测):摘出来并入别名流程,不落卡体字段
alias_new = list(alias_new or []) + \
[a for a in (changes.pop("别名新增", None) or []) if isinstance(a, str)]
for k, v in changes.items():
if valid_keys is not None and k not in valid_keys:
# 越合同字段:裁剪留审计(方案守卫条款),畸形长 key(条目误当 key)一并挡下
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, ("越合同:" + str(k))[:100],
json.dumps(v, ensure_ascii=False)[:2000], TENANT))
continue
if k == "一句话摘要":
# 摘要住卡顶层而非"字段"子对象——此前写进子对象成"影子摘要"且顶层摘要
# 从不更新(抽检 M4:600 章前的旧摘要一直挂着)。特判写顶层,同守水位闸。
if isinstance(v, str) and v.strip() and not _is_garbage(v):
if win_no < watermark:
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, "迟到覆写弃用:一句话摘要",
json.dumps(v, ensure_ascii=False)[:2000], TENANT))
else:
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,%s,%s,%s)""",
(draft_id, win_no, "一句话摘要",
json.dumps(payload.get("一句话摘要"), ensure_ascii=False),
json.dumps(v.strip(), ensure_ascii=False), TENANT))
payload["一句话摘要"] = v.strip()
continue
if k in MILESTONE_FIELDS:
# 里程碑对象数组(演变历程):走对象合并路径——去重按台阶内容、排序按真实章号,
# 不走下面处理字符串条目([窗N] 前缀)的老路径。迟到窗追加无害:对象自带章号,
# 乱序由 _chapter_sort_key 排序纠正,故不设水位闸(与字符串追加同策略)。
base = fields.get(k) if isinstance(fields.get(k), list) else ([fields[k]] if fields.get(k) else [])
merged, rejected = _merge_milestones(base, v, win_no)
for rj in rejected:
# 垃圾/空台阶里程碑拒收留审计(对齐字符串路径的垃圾拦截守卫)
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, "垃圾拦截:演变历程",
json.dumps(rj, ensure_ascii=False)[:2000], TENANT))
fields[k] = merged
continue
is_append = k in APPEND_FIELDS or isinstance(v, list) or isinstance(fields.get(k), list)
# 覆写值以增量口吻开头=模型把增量当全量(抽检#4:新值会抹掉旧基线),转追加不覆写
if not is_append and isinstance(v, str) and \
any(v.lstrip().startswith(w) for w in ("新增", "本窗", "另外", "此外")):
is_append = True
if is_append:
items = v if isinstance(v, list) else [v]
old = fields.get(k) if isinstance(fields.get(k), list) else ([fields[k]] if fields.get(k) else [])
seen = {_strip_prefix(x) for x in old} # 同文去重(抽检#1 重复病)
for it in items:
# 巨型粘连拆分(唐灵窗51 实测:模型把整个历史弧线连成一条"…;[窗2]…;[窗3]…"
# 输出成单条目)——按";[窗N]"边界拆段,带原窗号的段保留原窗号,其余记本窗
for seg in re.split(r";\s*(?=\[[窗本])", str(it)):
m0 = re.match(r"^\[窗(\d+)\]\s*(.*)", seg.strip(), flags=re.S)
seg_win, body = (int(m0.group(1)), m0.group(2)) if m0 else (win_no, seg)
core = _strip_tail(_strip_prefix(body)) # 剥模型自带前缀+尾部拼接残渣
if core and _is_garbage(core):
# 结构垃圾条目拒收留审计(窗75 事故根治:变更包字符串化混进字段值)
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, ("垃圾拦截:" + str(k))[:100],
str(core)[:2000], TENANT))
continue
if core and core not in seen:
old.append(f"[窗{seg_win}] {core}")
seen.add(core)
# 窗序归位(抽检 M2:补跑/重试窗条目尾插致时间线倒流)——稳定排序,同窗保持原序
fields[k] = sorted(old, key=_entry_win)
elif isinstance(v, str) and _is_garbage(v):
# 结构垃圾覆写值拒收留审计(同窗75 事故根治)
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, ("垃圾拦截:" + str(k))[:100], v[:2000], TENANT))
elif win_no < watermark:
# 迟到覆写弃用:本窗时序早于卡已生长到的窗位,覆写会让卡态倒流——
# 只留审计(标记可查),卡体保持高窗态不动
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(draft_id, win_no, ("迟到覆写弃用:" + str(k))[:100],
json.dumps(v, ensure_ascii=False)[:2000], TENANT))
else:
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,%s,%s,%s)""",
(draft_id, win_no, k, json.dumps(fields.get(k), ensure_ascii=False)
if fields.get(k) is not None else None,
json.dumps(v, ensure_ascii=False), TENANT))
# 覆写值同样剥尾部拼接残渣(批9 样张走查同款病灶,覆写路一并守住)
fields[k] = _strip_tail(v) if isinstance(v, str) else v
for al0 in alias_new or []:
al = _clean_alias(al0)
if al and al != payload.get("名称"):
payload.setdefault("别名", [])
if al not in payload["别名"]:
payload["别名"].append(al)
conn.execute(
"""INSERT INTO example_upgrade_alias
(work_id, canonical_name, alias, evidence_window, verdict_by, tenant_id)
VALUES (%s,%s,%s,%s,'ai',%s)
ON CONFLICT (tenant_id, work_id, alias) DO NOTHING""",
(payload.get("_work_id") or 0, payload.get("名称"), al, win_no, TENANT))
conn.execute(
"""UPDATE muse_knowledge_draft SET draft_payload=%s, revision=revision+1,
updater='upgrade' WHERE id=%s""",
(json.dumps(payload, ensure_ascii=False), draft_id))
conn.execute(
"""INSERT INTO example_upgrade_card_state (draft_id, work_id, watermark_window, tenant_id)
VALUES (%s,%s,%s,%s)
ON CONFLICT (draft_id) DO UPDATE
SET watermark_window=GREATEST(example_upgrade_card_state.watermark_window,
EXCLUDED.watermark_window),
update_time=now()""",
(draft_id, payload.get("_work_id") or 0, win_no, TENANT))
def new_card(conn, work_id, win_no, ent):
"""立初卡:payload 全字段以旧值=NULL 入审计(G7,错认拆回可还原初始态)。
名称剥括号注(抽检#7 根因):「白色游魂(无名侦察兵)」这类名称使预扫精确匹配失明
——主名之外的括号内容若像名字则转别名,否则丢弃。"""
import re as _re
raw = ent["名称"].strip()
m = _re.match(r"^(.+?)[((](.+?)[))]\s*$", raw)
extra_alias = []
if m:
raw = m.group(1).strip()
note = _clean_alias(m.group(2))
if note:
extra_alias.append(note)
# 初卡字段过守卫(深空 4917 实测:立卡路不走 merge_card,粘连/自造前缀/垃圾
# 原样入库——更新路守了、立卡路漏了):列表值逐条拆分、剥前缀、垃圾拦截、带窗号
fields0 = {}
for k, v in (ent.get("字段", {}) or {}).items():
if k in MILESTONE_FIELDS:
# 里程碑字段(演变历程):初卡即走对象合并(登场/首个进化台阶),去重排序;
# 初卡尚无 draft_id 无法留审计,垃圾条目直接过滤(与下方字符串路径初卡同策略)。
merged, _ = _merge_milestones([], v, win_no)
fields0[k] = merged
continue
if not isinstance(v, list):
fields0[k] = v
continue
out, seen = [], set()
for it in v:
for seg in re.split(r";\s*(?=\[[窗本])", str(it)):
core = _strip_tail(_strip_prefix(seg))
if not core or _is_garbage(core) or core in seen:
continue
out.append(f"[窗{win_no}] {core}")
seen.add(core)
fields0[k] = out
payload = {"type": ent["型"], "名称": raw,
"别名": [x for x in ([_clean_alias(a) for a in ent.get("别名", [])] + extra_alias) if x],
"一句话摘要": ent.get("一句话摘要", ""),
"字段": fields0,
"出场章": sorted(set(ent.get("出场章", []))),
"来源": f"升格@窗{win_no}", "状态": "草稿",
"目标库": "本书作品库", "可见范围": "本书私有",
"_work_id": work_id}
did = 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""",
(work_id, json.dumps(payload, ensure_ascii=False), SOURCE_TYPE, work_id,
TENANT)).fetchone()[0]
for k, v in payload["字段"].items():
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(did, win_no, k, json.dumps(v, ensure_ascii=False), TENANT))
# 初始别名同步入 alias 表(抽检 L3 根源:立卡别名只写卡内不入表,
# 别名表覆盖率仅 26%,"别名精确判重"层大面积空转导致同人两卡漏并)
for al in payload["别名"]:
conn.execute(
"""INSERT INTO example_upgrade_alias
(work_id, canonical_name, alias, evidence_window, verdict_by, tenant_id)
VALUES (%s,%s,%s,%s,'init',%s)
ON CONFLICT (tenant_id, work_id, alias) DO NOTHING""",
(work_id, raw, al, win_no, TENANT))
conn.execute(
"""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
def undo_window(conn, work_id, win_no):
"""同窗重跑先撤销(防重跑自噬):按审计还原覆写、按窗号删追加条目与留档。"""
# 还原覆写字段(倒序还原,先写的最后还原到最初旧值)
rows = conn.execute(
"""SELECT a.draft_id, a.field_name, a.old_value FROM example_upgrade_audit a
JOIN muse_knowledge_draft d ON d.id=a.draft_id
WHERE a.tenant_id=%s AND d.work_id=%s AND a.window_no=%s ORDER BY a.id DESC""",
(TENANT, work_id, win_no)).fetchall()
for did, fname, old in rows:
payload = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
if old is None:
payload.get("字段", {}).pop(fname, None)
else:
payload.setdefault("字段", {})[fname] = json.loads(old)
conn.execute("UPDATE muse_knowledge_draft SET draft_payload=%s WHERE id=%s",
(json.dumps(payload, ensure_ascii=False), did))
conn.execute("""DELETE FROM example_upgrade_audit WHERE tenant_id=%s AND window_no=%s
AND draft_id IN (SELECT id FROM muse_knowledge_draft WHERE work_id=%s)""",
(TENANT, win_no, work_id))
# 删本窗追加条目([窗N] 前缀)与初立于本窗的卡(audit 已删,靠 card_state 水位判初窗不可靠,
# 初卡以「来源=升格@窗N」标记识别)
for did, payload in conn.execute(
"""SELECT id, draft_payload FROM muse_knowledge_draft
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE""",
(TENANT, work_id, SOURCE_TYPE)).fetchall():
if payload.get("来源") == f"升格@窗{win_no}":
conn.execute("UPDATE muse_knowledge_draft SET deleted=TRUE WHERE id=%s", (did,))
# 同步清卡水位行(fable 复验实证:深空回滚删 400+ 初卡后 card_state 僵尸行
# 残留,按 draft JOIN 不过滤 deleted 的查询会捞出僵尸卡)
conn.execute("DELETE FROM example_upgrade_card_state WHERE draft_id=%s", (did,))
continue
tag, changed = f"[窗{win_no}] ", False
for k, v in list(payload.get("字段", {}).items()):
if isinstance(v, list):
# 字符串条目按 [窗N] 前缀删;里程碑对象(演变历程)按内部 _win 溯源键删——
# 对象无窗号前缀,撤销靠 _win 精确识别本窗新增,否则同窗重跑会 double-append。
nv = [x for x in v
if not (isinstance(x, str) and x.startswith(tag))
and not (isinstance(x, dict) and x.get("_win") == win_no)]
if len(nv) != len(v):
payload["字段"][k], changed = nv, True
if changed:
conn.execute("UPDATE muse_knowledge_draft SET draft_payload=%s WHERE id=%s",
(json.dumps(payload, ensure_ascii=False), did))
conn.execute("DELETE FROM example_upgrade_presence WHERE tenant_id=%s AND work_id=%s AND window_no=%s",
(TENANT, work_id, win_no))
conn.execute("""DELETE FROM example_upgrade_alias WHERE tenant_id=%s AND work_id=%s
AND evidence_window=%s""", (TENANT, work_id, win_no))
conn.commit()
# ── 语义判重(P1,设计稿 §8.2):打开 v6 已设计、暂时关着的「嵌入近邻 + M3 终判」那级 ──
# 治病根 4:机械判重只比名字字符串,改名("影杀者"→"IV代纯机械机甲·影杀者")/跨型指代就漏并。
# 默认关(--semantic-dedup 开启):语义召回依赖同书卡已 embed 落库;试跑期嵌入延后,无向量时优雅空转。
def _entity_embed_text(ent):
"""判重用嵌入文本:【型】名称:摘要 + 关键字段摘选(与检索端 embed 文本语义对齐)。
兼容新实体 ent(键:型/名称/一句话摘要/字段);下划线内部键(如 _win)不进嵌入文本。"""
t = ent.get("型") or ent.get("type") or ""
fields = ent.get("字段") or {}
body = "\n".join(f"{k}:{v}" for k, v in fields.items()
if v and k not in ("名称", "一句话摘要") and not str(k).startswith("_"))
return f"【{t}】{ent.get('名称', '')}:{ent.get('一句话摘要', '')}\n{body}"[:4000]
def recall_neighbors(conn, sess, work_id, ent, top=6):
"""算候选向量→在 example_knowledge_embedding 召回同书 ≥阈值 近邻卡。
返回 [(did, 型, 名称, 摘要, 相似度)…] 按相似度降序。
连接纪律:只读短查询,沿用窗内 conn(与 observe/update 同模式,keepalives 兜底)。
嵌入表无同书向量(试跑期未 embed)时返回 []——优雅降级,判重回退纯机械,绝不阻断主流程。"""
vecs, bad = embed_texts(sess, [_entity_embed_text(ent)])
if bad or not vecs or vecs[0] is None:
return []
qvec = json.dumps(vecs[0])
rows = conn.execute(
"""SELECT d.id, d.draft_payload->>'type', d.draft_payload->>'名称',
d.draft_payload->>'一句话摘要', 1 - (e.embedding <=> %s::vector) AS score
FROM example_knowledge_embedding e
JOIN muse_knowledge_draft d ON d.id = e.draft_id
WHERE e.tenant_id=%s AND e.deleted=FALSE AND d.deleted=FALSE
AND d.source_type=%s AND d.work_id=%s
ORDER BY score DESC LIMIT %s""",
(qvec, TENANT, SOURCE_TYPE, work_id, top)).fetchall()
return [(r[0], r[1], r[2], r[3], float(r[4])) for r in rows
if float(r[4]) >= DEDUP_SIM_THRESHOLD]
def dedup_judge_prompt(ent, neighbors):
"""M3 终判 prompt:候选新实体 vs 每个同书近邻,判 同一实体 / 前身 / 后继 / 无关。"""
nb = "\n".join(f"{i + 1}. 卡号{did}|{t}|{nm}|{brief or ''}"
for i, (did, t, nm, brief, _s) in enumerate(neighbors))
return f"""【功能指令(parse-book 作品面升格·语义判重终判)】
下面是一个"候选新实体"和若干"同书既有卡"(向量召回的近邻)。逐一判断候选与每张近邻卡的关系,四选一:
- 同一实体:同一对象的改名/化名/不同侧面(如"影杀者"与"IV代纯机械机甲·影杀者")——**仅同型可判**。
- 前身:候选是该近邻卡的上一代/来源(同一条进化链的相邻代际,如"铁头(一代)"之于"铁卫(二代)")。
- 后继:候选是该近邻卡的下一代/继承者。
- 无关:只是题材相近,各自独立。
判据:看名称/摘要/字段是否指向同一对象或同一条演变链;**跨型(如具体机甲 item vs 整套体系 power_system)绝不判同一实体,最多判前身/后继**——一整套体系不等于其中一台机体。
【候选新实体】
型={ent.get("型", "")}|名称={ent.get("名称", "")}|摘要={ent.get("一句话摘要", "")}
字段:{json.dumps(ent.get("字段", {}), ensure_ascii=False)[:600]}
【同书近邻卡】
{nb}
【输出规则(只输出一个 JSON 对象)】
{{"判定": [{{"卡号": 数字, "关系": "同一实体|前身|后继|无关"}}]}}"""
def semantic_dedup(conn, sess, work_id, ent):
"""语义判重裁决。返回 (verdict, data):
('merge', (did, 近邻名称)) 同型·同一实体 → 并卡(治改名漏并)
('chain', [(did, 关系, 名称)]) 前身后继 → 不并卡但记串链候选
('new', None) 无近邻或全判无关 → 各自立卡
同型才允许 merge;跨型即便 M3 判同一实体也降级为串链候选(设计稿 §8.2:跨型仅提示、不自动并)。
判重是增益非必需:召回失败/无向量/终判异常一律保守返回 new(不并可后补,误并难回退)。"""
neighbors = recall_neighbors(conn, sess, work_id, ent)
if not neighbors:
return "new", None
try:
data, _ = m3_json(dedup_judge_prompt(ent, neighbors), "MiniMax-M3", ("判定",))
except (SensitiveHardStop, RuntimeError):
return "new", None
nb_type = {did: t for did, t, _, _, _ in neighbors}
nb_name = {did: nm for did, _, nm, _, _ in neighbors}
ent_type = ent.get("型", "")
chain = []
for j in [x for x in (data.get("判定") or []) if isinstance(x, dict)]:
did, rel = j.get("卡号"), j.get("关系")
if did not in nb_type:
continue
same_type = nb_type[did] == ent_type
if rel == "同一实体" and same_type:
return "merge", (did, nb_name[did]) # 近邻按相似度降序,第一个同型同一实体即采
if rel in ("前身", "后继"):
chain.append((did, rel, nb_name[did]))
elif rel == "同一实体": # 跨型判同一实体不可信 → 降级串链候选
chain.append((did, "前身后继待定", nb_name[did]))
return ("chain", chain) if chain else ("new", None)
# ── 命令 ──
@click.group()
def cli():
"""作品面升格(正文窗直抽 · 全实体统一生长)"""
@cli.command()
@click.option("--work-id", type=int, required=True)
def windows(work_id):
"""机械切正文窗(幂等)。"""
with psycopg.connect(DSN) as conn:
total, new = cut_windows(conn, work_id)
click.echo(f"work={work_id} 切窗完成:全书 {total} 窗(本次新建 {new} 行)")
@cli.command()
@click.option("--work-id", type=int, required=True)
@click.option("--max-windows", type=int, default=0, help="本次最多跑几个窗(0=不限)")
@click.option("--max-calls", type=int, default=0, help="本次 LLM 调用上限(0=不限,含敏感失败)")
@click.option("--model", default="MiniMax-M3", show_default=True)
@click.option("--redo-window", type=int, default=0, help="指定窗号强制重跑(先撤销后重写)")
@click.option("--semantic-dedup", "semantic_on", is_flag=True,
help="开启语义判重(P1):立卡前召回同书近邻+M3终判治改名/跨型漏并"
"(需同书已 embed 落库;默认关=试跑期嵌入延后)")
def run(work_id, max_windows, max_calls, model, redo_window, semantic_on):
"""按窗顺序跑升格:断点续跑跳过 done 窗;敏感硬停=窗 failed+书停。"""
calls = {"n": 0} # 调用计数(含敏感失败换模型的次数由 m3_json 内部消化,此处计成功轮次)
# 语义判重嵌入会话(仅开启时建;禁系统代理,走内网直连)
embed_sess = _embed_session() if semantic_on else None
def call(prompt, need_keys):
calls["n"] += 1
return m3_json(prompt, model, need_keys, system=IDENTITY)
with psycopg.connect(DSN) as conn:
title = conn.execute("SELECT title FROM muse_content_work WHERE id=%s",
(work_id,)).fetchone()[0]
contracts = load_entity_contracts(conn)
if redo_window:
row = conn.execute(
"""SELECT window_no, from_chapter, to_chapter FROM example_upgrade_window
WHERE tenant_id=%s AND work_id=%s AND window_no=%s AND deleted=FALSE""",
(TENANT, work_id, redo_window)).fetchone()
wins = [row] if row else []
if wins:
click.echo(f"[撤销] 窗{redo_window} 旧写入回滚中…")
undo_window(conn, work_id, redo_window)
else:
wins = conn.execute(
"""SELECT window_no, from_chapter, to_chapter, status FROM example_upgrade_window
WHERE tenant_id=%s AND work_id=%s AND status!='done' AND deleted=FALSE
ORDER BY from_chapter""", (TENANT, work_id)).fetchall()
wins = [(r[0], r[1], r[2]) for r in wins]
# failed 窗重跑前必须先撤销旧写入(部分写入直接重跑会 double-append 追加字段)
for r0 in conn.execute(
"""SELECT window_no FROM example_upgrade_window
WHERE tenant_id=%s AND work_id=%s AND status='failed' AND deleted=FALSE""",
(TENANT, work_id)).fetchall():
click.echo(f"[撤销] failed 窗{r0[0]} 旧写入回滚后重跑")
undo_window(conn, work_id, r0[0])
done_n = 0
# 串行铁律(创始人 2026-07-15 拍板):窗与窗是串行生长——后窗的预扫/判重/更新
# 全依赖前窗长成的卡。窗失败绝不跳窗(跳窗=知识断层+事后补跑有覆盖风险),
# 而是当场撤销半写入→整窗重试一次(挡偶发病);仍败→停书,断点就在本窗,
# 下次启动从本窗续跑(failed 先撤销机制),串行语义天然无损。
wi = 0
retried = False # 当前窗是否已当场重试过
while wi < len(wins):
win_no, a, b = wins[wi]
if max_windows and done_n >= max_windows:
break
if max_calls and calls["n"] >= max_calls:
click.echo(f"⏸ 调用闸 {max_calls} 已到,停在窗{win_no} 之前")
break
try:
with psycopg.connect(DSN) as conn:
text = load_window_text(conn, work_id, a, b)
name_map, presence = load_known(conn, work_id)
onstage = prescan(name_map, text)
# ① 观察
obs, usage = call(observe_prompt(contracts, title, a, b, text, onstage),
("新名字", "已知实体新信息", "纯出场"))
# 模型输出防御(深空窗5实测:列表元素偶为裸字符串,.get 直接炸)——统一只留 dict 元素
for k in ("新名字", "已知实体新信息", "纯出场"):
obs[k] = [x for x in (obs.get(k) or []) if isinstance(x, dict)]
# ② 判重+立卡门槛(机械)
to_update = {} # draft_id -> [观察点…]
for ent in obs.get("新名字", []):
nm = (ent.get("名称") or "").strip()
if not nm:
continue
hint = (ent.get("疑似别名指向") or "").strip()
if nm in name_map: # 观察漏看在场清单:直接归并
did = name_map[nm][0]
to_update.setdefault(did, []).append(
f"(新名字归并){ent.get('一句话摘要', '')}")
continue
if hint and hint in name_map: # 疑似别名:初卡转观察材料(G3)
did = name_map[hint][0]
to_update.setdefault(did, []).append(
f"(别名「{nm}」并入)初卡材料:{json.dumps(ent, ensure_ascii=False)[:400]}")
conn.execute(
"""INSERT INTO example_upgrade_alias
(work_id, canonical_name, alias, evidence_window, verdict_by, tenant_id)
VALUES (%s,%s,%s,%s,'ai',%s)
ON CONFLICT (tenant_id, work_id, alias) DO NOTHING""",
(work_id, hint, nm, win_no, TENANT))
continue
# 同型名称互为子串(抽检 H5:「果子」vs「开心果子」同人两卡):
# 不直接立卡,转观察材料并入既有卡由更新步 AI 甄别(软防护,留复核标记)
sub_hit = next(
(ex for ex, (d0, t0, _) in name_map.items()
if t0 == ent.get("型", "") and len(nm) >= 2 and len(ex) >= 2
and nm != ex and (nm in ex or ex in nm)), None)
if sub_hit:
did = name_map[sub_hit][0]
to_update.setdefault(did, []).append(
f"(名称疑似同一实体「{nm}」≈「{sub_hit}」,请甄别后再并入)"
f"初卡材料:{json.dumps(ent, ensure_ascii=False)[:400]}")
conn.execute(
"""INSERT INTO example_upgrade_alias
(work_id, canonical_name, alias, evidence_window, verdict_by, tenant_id)
VALUES (%s,%s,%s,%s,'substr',%s)
ON CONFLICT (tenant_id, work_id, alias) DO NOTHING""",
(work_id, sub_hit, nm, win_no, TENANT))
continue
chaps = set(ent.get("出场章", []))
hist = presence.get((ent.get("型", ""), nm), set())
if len(chaps | hist) >= 2: # 跨章(含跨窗合计)→ 立卡
if hist: # 用留档补足初卡出场章
ent["出场章"] = sorted(chaps | hist)
# 语义判重(P1,--semantic-dedup 开启且同书已 embed 时生效):机械判重
# (名字/别名/子串)之后、立卡之前,召回同书近邻交 M3 终判治改名/跨型漏并。
if semantic_on:
verdict, vd = semantic_dedup(conn, embed_sess, work_id, ent)
if verdict == "merge": # 同型同一实体:并入既有卡,不另立
did0, canon = vd
to_update.setdefault(did0, []).append(
f"(语义判重·「{nm}」并入同一实体)初卡材料:"
f"{json.dumps(ent, ensure_ascii=False)[:400]}")
conn.execute(
"""INSERT INTO example_upgrade_alias
(work_id, canonical_name, alias, evidence_window, verdict_by, tenant_id)
VALUES (%s,%s,%s,%s,'semantic',%s)
ON CONFLICT (tenant_id, work_id, alias) DO NOTHING""",
(work_id, canon, nm, win_no, TENANT))
continue
if verdict == "chain": # 前身后继:仍立卡,串链关系记候选审计
did = new_card(conn, work_id, win_no, ent)
name_map[nm] = (did, ent.get("型", ""), ent.get("一句话摘要", ""))
# 只记候选提示、不自动写「前身/后继」字段——避免误串,链接由人工/后续确认落字段
for did2, rel, nm2 in vd:
conn.execute(
"""INSERT INTO example_upgrade_audit
(draft_id, window_no, field_name, old_value, new_value, tenant_id)
VALUES (%s,%s,%s,NULL,%s,%s)""",
(did, win_no, ("串链候选:" + str(rel))[:100],
json.dumps({"对方卡号": did2, "对方名称": nm2},
ensure_ascii=False), TENANT))
continue
did = new_card(conn, work_id, win_no, ent)
name_map[nm] = (did, ent.get("型", ""), ent.get("一句话摘要", ""))
else: # 单章龙套 → 留档(G4)
for ch in (chaps or {a}):
conn.execute(
"""INSERT INTO example_upgrade_presence
(work_id, window_no, chapter_no, entity_type, name,
observation, tenant_id)
VALUES (%s,%s,%s,%s,%s,%s,%s)""",
(work_id, win_no, ch, ent.get("型", ""), nm,
ent.get("一句话摘要", ""), TENANT))
pres_add = {} # did -> set(出场章)——被更新的卡也要记出场章(抽检#2:只走纯出场路径整窗丢章)
for it in obs.get("已知实体新信息", []):
nm = (it.get("名称") or "").strip()
if nm in name_map and it.get("观察点"):
did = name_map[nm][0]
to_update.setdefault(did, []).append(it["观察点"])
pres_add.setdefault(did, set()).update(
c for c in (it.get("出场章") or []) if isinstance(c, int))
# ④ 卡更新(分批≤6)
items = sorted(to_update.items())
for i in range(0, len(items), UPDATE_BATCH):
batch = items[i:i + UPDATE_BATCH]
cards = []
for did, obs_pts in batch:
p = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
cards.append((did, p, obs_pts))
# 缺"更新"键宽容为空批:prompt 教"无变化的卡不输出",某批恰好全无变化时
# 模型会顺势连键一起省(批7实测 6 窗全死于此)。语义上缺键≈空批,按空批放行
# 并留警告日志可审计;其余格式错误(乱码/解析失败)仍原样抛、窗照 fail。
try:
upd, _ = call(update_prompt(contracts, title, a, b, text, cards), ("更新",))
except RuntimeError as e:
if "缺少必需键" not in str(e):
raise
print(f"[宽容] 窗{win_no} 更新批缺键按空批放行: {str(e)[:80]}", file=sys.stderr)
upd = {"更新": []}
valid = {d for d, _, _ in cards}
# 每卡按其型的合同 key 集校验(+一句话摘要),越合同 key 裁剪留审计
did2keys = {d: {f["key"] for f in contracts.get(p.get("type"), {}).get("字段", [])}
| {"一句话摘要"} for d, p, _ in cards}
for u in [x for x in (upd.get("更新") or []) if isinstance(x, dict)]:
if u.get("draft_id") in valid and (u.get("变更字段") or u.get("别名新增")):
merge_card(conn, u["draft_id"], win_no,
u.get("变更字段"), u.get("别名新增"),
valid_keys=did2keys.get(u["draft_id"]))
# ⑤ 关系增量(核心角色=本窗有更新的 character + 在场 character,≤8)
char_cards = []
seen = set()
for did in list(to_update.keys()):
p = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
if p.get("type") == "character" and did not in seen:
char_cards.append((did, p))
seen.add(did)
for nm, (did, t, _) in onstage.items():
if t == "character" and did and did not in seen and len(char_cards) < 8:
p = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
char_cards.append((did, p))
seen.add(did)
char_cards = char_cards[:8]
if len(char_cards) >= 2:
rels = conn.execute(
"""SELECT id, draft_payload FROM muse_knowledge_draft
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE
AND draft_payload->>'type'=%s""",
(TENANT, work_id, SOURCE_TYPE, RELATION_TYPE)).fetchall()
rel_out, _ = call(relation_prompt(contracts, title, a, b, text,
char_cards, rels), ("关系",))
id2name = {d: p.get("名称") for d, p in char_cards}
exist = {tuple(sorted((r.get("甲方draft"), r.get("乙方draft")))): (rid, r)
for rid, r in rels
if r.get("甲方draft") and r.get("乙方draft")}
for r in [x for x in (rel_out.get("关系") or []) if isinstance(x, dict)]:
ja, yi = r.get("甲方"), r.get("乙方")
if ja not in id2name or yi not in id2name or ja == yi:
continue
key = tuple(sorted((ja, yi)))
evo_core = _strip_prefix(r.get("本窗演变", ""))
if key in exist:
# 落点统一(抽检#3):其他字段与演变轨迹全并进 rp["字段"],
# 顶层遗留的演变轨迹迁移进字段后删除,避免双份分裂
rid, rp = exist[key]
f2 = rp.setdefault("字段", {})
if "演变轨迹" in rp: # 存量顶层迁移
legacy = rp.pop("演变轨迹")
base = f2.get("演变轨迹") or []
f2["演变轨迹"] = (base if isinstance(base, list) else [base]) + \
(legacy if isinstance(legacy, list) else [legacy])
for k2, v2 in (r.get("其他字段") or {}).items():
if k2 == "演变轨迹":
continue # 演变只走本窗演变通道
f2[k2] = v2 # 覆写(当前状态等保最新)
if evo_core:
lst = f2.setdefault("演变轨迹", [])
if not isinstance(lst, list):
lst = [lst]
f2["演变轨迹"] = lst
if evo_core not in {_strip_prefix(x) for x in lst}:
lst.append(f"[窗{win_no}] {evo_core}")
if r.get("关系类型"):
rp["关系类型"] = r["关系类型"]
conn.execute(
"""UPDATE muse_knowledge_draft SET draft_payload=%s,
revision=revision+1 WHERE id=%s""",
(json.dumps(rp, ensure_ascii=False), rid))
else:
f2 = dict(r.get("其他字段") or {})
f2["演变轨迹"] = [f"[窗{win_no}] {evo_core}"] if evo_core else []
rp = {"type": RELATION_TYPE,
"名称": f"{id2name[ja]}×{id2name[yi]}",
"甲方draft": ja, "乙方draft": yi,
"甲方名称": id2name[ja], "乙方名称": id2name[yi],
"关系类型": r.get("关系类型", ""),
"字段": f2,
"来源": f"升格@窗{win_no}", "状态": "草稿",
"目标库": "本书作品库", "可见范围": "本书私有",
"_work_id": work_id}
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)""",
(work_id, json.dumps(rp, ensure_ascii=False), SOURCE_TYPE,
work_id, TENANT))
# ⑥ 机械收尾:出场章统一并入(纯出场 + 被更新卡,抽检#2)+ 窗置 done
for it in obs.get("纯出场", []):
nm = (it.get("名称") or "").strip()
if nm in name_map:
pres_add.setdefault(name_map[nm][0], set()).update(
c for c in (it.get("出场章") or []) if isinstance(c, int))
for did, chs in pres_add.items():
if not chs:
continue
p = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
p["出场章"] = sorted(set(p.get("出场章", [])) | chs)
conn.execute("UPDATE muse_knowledge_draft SET draft_payload=%s WHERE id=%s",
(json.dumps(p, ensure_ascii=False), did))
conn.execute(
"""UPDATE example_upgrade_window SET status='done', error_message=NULL,
updater='upgrade' WHERE tenant_id=%s AND work_id=%s AND window_no=%s""",
(TENANT, work_id, win_no))
conn.commit()
new_n = len([e for e in obs.get('新名字', [])])
click.echo(f" 窗{win_no}✓ ({a}-{b}章) 在场{len(onstage)} 新名字{new_n} "
f"更新卡{len(to_update)} 调用累计{calls['n']}")
done_n += 1
wi += 1
retried = False
except SensitiveHardStop as e:
with psycopg.connect(DSN) as conn:
conn.execute(
"""UPDATE example_upgrade_window SET status='failed', error_message=%s
WHERE tenant_id=%s AND work_id=%s AND window_no=%s""",
(str(e)[:500], TENANT, work_id, win_no))
conn.commit()
click.echo(f" ⛔ 窗{win_no} 敏感降级链全失败,本书升格硬停:{e}")
return
except Exception as e:
with psycopg.connect(DSN) as conn:
conn.execute(
"""UPDATE example_upgrade_window SET status='failed', error_message=%s
WHERE tenant_id=%s AND work_id=%s AND window_no=%s""",
(str(e)[:500], TENANT, work_id, win_no))
conn.commit()
if not retried:
# 串行铁律第一层:当场撤销半写入→同窗立即重试(挡网络瞬断/偶发格式病)
with psycopg.connect(DSN) as conn:
undo_window(conn, work_id, win_no)
conn.commit()
retried = True
click.echo(f" ↻ 窗{win_no} 失败,撤销后当场重试(串行铁律不跳窗): {str(e)[:150]}")
continue
# 第二层重试仍败:停书,断点=本窗;下次启动 failed 先撤销、从本窗续跑
click.echo(f" ⛔ 窗{win_no} 当场重试仍失败,本书升格停(断点=本窗,续跑从此接续): {str(e)[:200]}")
return
click.echo(f"《{title}》本次完成 {done_n} 窗,LLM 调用 {calls['n']} 次")
@cli.command()
@click.option("--work-id", type=int, required=True)
def status(work_id):
"""升格进度:窗状态/卡数/留档数。"""
with psycopg.connect(DSN) as conn:
title = conn.execute("SELECT title FROM muse_content_work WHERE id=%s",
(work_id,)).fetchone()[0]
w = conn.execute(
"""SELECT count(*) FILTER (WHERE status='done'), count(*) FILTER (WHERE status='failed'),
count(*) FROM example_upgrade_window
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""",
(TENANT, work_id)).fetchone()
cards = conn.execute(
"""SELECT draft_payload->>'type', count(*) FROM muse_knowledge_draft
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE
GROUP BY 1 ORDER BY 2 DESC""", (TENANT, work_id, SOURCE_TYPE)).fetchall()
pres = conn.execute(
"SELECT count(*) FROM example_upgrade_presence WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE",
(TENANT, work_id)).fetchone()[0]
ali = conn.execute(
"SELECT count(*) FROM example_upgrade_alias WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE",
(TENANT, work_id)).fetchone()[0]
click.echo(f"《{title}》窗 {w[0]}done/{w[1]}failed/{w[2]}total | "
f"卡 {', '.join(f'{t}:{n}' for t, n in cards) or '0'} | 留档{pres} 别名{ali}")
if __name__ == "__main__":
cli()