框架: 作品面升格管线parse_upgrade.py落地(v6统一生长机制,机动窗1-2实测验证)

This commit is contained in:
zizi 2026-07-14 22:12:50 +08:00
parent 4369984c9c
commit b7eb71b370

View File

@ -0,0 +1,615 @@
#!/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 作品面升格·卡增量更新)】
下列每张卡给出当前卡全文本窗观察点对照本窗正文**只输出需要变更的字段**
- 覆写类字段性格底色/说话方式/当前状态等标量输出更准确的新值整字段替换
- 追加类字段成长弧线/演变轨迹等数组只输出本窗新增条目不要重抄旧条目
- 没有变化的字段不要输出整卡无实质变化则不输出该卡
纪律以正文为证据不脑补字段 key 必须来自该型合同每卡别名有新发现可在别名新增里给只收**该实体自己**的新别名他人对它的称呼算它对别的实体的称呼不算窗2实测魔鬼金是主角喊宠物机的诨名错报成了主角别名
相关型字段合同
{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}"""
# ── 合并与审计(三类字段演进 + 撤销依据)──
def merge_card(conn, draft_id, win_no, changes, alias_new):
"""按 5.1 三类规则合并变更字段:数组/白名单=追加(条目带窗号),标量=覆写留审计。"""
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():
is_append = k in APPEND_FIELDS or isinstance(v, list) or isinstance(fields.get(k), list)
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 [])
fields[k] = old + [f"[窗{win_no}] {it}" for it in items if it]
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 al in alias_new or []:
if al and al.strip() and al.strip() != payload.get("名称"):
payload.setdefault("别名", [])
if al.strip() not in payload["别名"]:
payload["别名"].append(al.strip())
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.strip(), 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错认拆回可还原初始态"""
payload = {"type": ent[""], "名称": ent["名称"].strip(),
"别名": [a.strip() for a in ent.get("别名", []) if a and a.strip()],
"一句话摘要": 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 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()
done_n = 0
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),
("新名字", "已知实体新信息", "纯出场"))
# ② 判重+立卡门槛(机械)
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))
for it in obs.get("已知实体新信息", []):
nm = (it.get("名称") or "").strip()
if nm in name_map and it.get("观察点"):
to_update.setdefault(name_map[nm][0], []).append(it["观察点"])
# ④ 卡更新分批≤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}
for u in upd.get("更新", []):
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("别名新增"))
# ⑤ 关系增量(核心角色=本窗有更新的 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 rel_out.get("关系", []):
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 = f"[窗{win_no}] {r.get('本窗演变', '')}"
if key in exist:
rid, rp = exist[key]
rp.setdefault("演变轨迹", []).append(evo)
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:
rp = {"type": RELATION_TYPE,
"名称": f"{id2name[ja]}×{id2name[yi]}",
"甲方draft": ja, "乙方draft": yi,
"甲方名称": id2name[ja], "乙方名称": id2name[yi],
"关系类型": r.get("关系类型", ""),
"演变轨迹": [evo], "字段": r.get("其他字段", {}) or {},
"来源": 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))
# ⑥ 机械收尾:纯出场留档 + 窗置 done
for it in obs.get("纯出场", []):
nm = (it.get("名称") or "").strip()
if nm in name_map:
did = name_map[nm][0]
p = conn.execute("SELECT draft_payload FROM muse_knowledge_draft WHERE id=%s",
(did,)).fetchone()[0]
chs = sorted(set(p.get("出场章", [])) | set(it.get("出场章", [])))
p["出场章"] = 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
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()
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()