724 lines
41 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 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
# ── 窗切割参数(方案 §五B35 万字/窗、1015 章,取保守双闸防观察输出过载)──
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"
# 追加类字段白名单(值为数组的字段一律追加;这些字符串字段也强制追加、条目带窗号)
APPEND_FIELDS = {"成长弧线", "演变轨迹", "大事记", "经历"}
# ── 库内合同元数据驱动公理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) 纯出场:清单中实体本窗出现但无实质新信息的,只报名称+出场章。
纪律:一次性龙套(单章无名或仅路过)不报进新名字;实体判据与字段以合同为准,无证据不填;不脑补。数值(战力/指数/排名等)必须正文原样出现才可写,禁止推算或编造。组织改组/合并产生的新组织是**新实体**(走新名字),不是旧组织的别名。「疑似别名指向」只在确为同一实体改名/化名时填。
【六型字段合同】
{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 作品面升格·卡增量更新)】
下列每张卡给出「当前卡全文」与「本窗观察点」。对照本窗正文,**只输出需要变更的字段**
- 覆写类字段(性格底色/说话方式/当前状态等标量):**必须输出该字段完整的新全量值**——旧值里仍然成立的信息要保留进新值,禁止只写"新增…"式增量(那会把旧信息抹掉);
- 追加类字段(成长弧线/演变轨迹等数组):只输出本窗新增条目(不要重抄旧条目,不要自己加 [窗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实测模型自造 "[窗本窗]",仅数字版剥不掉)
WIN_PREFIX_RE = _re.compile(r"^(?:\[[窗本][^\]]{0,6}\]\s*)+")
return WIN_PREFIX_RE.sub("", str(s)).strip()
def _clean_alias(al):
"""别名机械准入(抽检#6拒收含括号注释/超长的备忘录式别名——预扫精确匹配永不命中=死数据。"""
al = (al or "").strip()
if not al or len(al) > 12:
return None
if any(c in al for c in "()。,,"):
return None
return al
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("字段", {})
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
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:
core = _strip_prefix(it) # 剥模型自带前缀(抽检#1 堆叠病)
if core and core not in seen:
old.append(f"[窗{win_no}] {core}")
seen.add(core)
fields[k] = old
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))
fields[k] = 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=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)
payload = {"type": ent[""], "名称": raw,
"别名": [x for x in ([_clean_alias(a) for a in ent.get("别名", [])] + extra_alias) if x],
"一句话摘要": ent.get("一句话摘要", ""),
"字段": ent.get("字段", {}) or {},
"出场章": 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))
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,))
continue
tag, changed = f"[窗{win_no}] ", False
for k, v in list(payload.get("字段", {}).items()):
if isinstance(v, list):
nv = [x for x in v if not (isinstance(x, str) and x.startswith(tag))]
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()
# ── 命令 ──
@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="指定窗号强制重跑(先撤销后重写)")
def run(work_id, max_windows, max_calls, model, redo_window):
"""按窗顺序跑升格:断点续跑跳过 done 窗;敏感硬停=窗 failed+书停。"""
calls = {"n": 0} # 调用计数(含敏感失败换模型的次数由 m3_json 内部消化,此处计成功轮次)
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
fail_streak = 0 # 连续失败计数单窗偶发失败网络瞬断continue连续 2 窗停书(系统性问题)
for win_no, a, b in wins:
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
chaps = set(ent.get("出场章", []))
hist = presence.get((ent.get("", ""), nm), set())
if len(chaps | hist) >= 2: # 跨章(含跨窗合计)→ 立卡
if hist: # 用留档补足初卡出场章
ent["出场章"] = sorted(chaps | hist)
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))
upd, _ = call(update_prompt(contracts, title, a, b, text, cards), ("更新",))
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
fail_streak = 0
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()
fail_streak += 1
# 偶发失败(网络瞬断/单窗输出畸形跳去下一窗failed 窗留给续跑(自动先撤销);
# 连续 2 窗失败=系统性问题,停书防连环烧额度
if fail_streak >= 2:
click.echo(f" ⛔ 连续 {fail_streak} 窗失败,本书升格停(末窗{win_no}: {str(e)[:200]}")
return
click.echo(f" ✗ 窗{win_no} 失败(已记录,续跑重试;连续第{fail_streak}次): {str(e)[:200]}")
continue
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()