576 lines
35 KiB
Python
576 lines
35 KiB
Python
#!/usr/bin/env python3
|
||
"""parse-book skill:M3 直调拆书执行器(B4-S3 重构:出卡权上移窗级)。
|
||
|
||
创始人拍板(2026-07-13):拆书内容生产 LLM=New-API MiniMax-M3(经 llm skill),
|
||
不再派 opus/haiku 子代理。B4 fable 审查裁决(2026-07-13):章级逐章出卡有三同根病
|
||
(同功 family 撞车/单章证不成跨章公式/间隔数字伪精确),治法=章级只产「范式候选线索」,
|
||
出卡在窗级聚类归并(复用大纲窗切分)——同一手法多章多次出现归并为一张母卡+实例章号,
|
||
间隔数字由实例章号差机械计算。
|
||
|
||
管线三步(顺序依赖):
|
||
chapters 逐章一次 M3:细纲+实体+范式候选线索 → example_parse_scaffold(正文只过一遍)
|
||
(parse_outline window:窗级大纲聚合 → example_parse_outline,窗=出卡窗的 SoT)
|
||
cards 逐窗一次 M3:窗内线索+细纲+阶段大纲 → 聚类出母卡 → parse_ingest cards 守卫入库
|
||
|
||
断点续跑:chapters 按 example_parse_task.scaffold_status 跳过;cards 按窗内是否已有活卡跳过。
|
||
"""
|
||
import json
|
||
import pathlib
|
||
import re
|
||
import subprocess
|
||
import sys
|
||
|
||
import click
|
||
import psycopg
|
||
|
||
# 统一走 llm skill 入口(trust_env/重试/<think>剥离/JSON 容错都在那边)
|
||
# 敏感/额度降级链已上收 llm.chat_governed(全局统一治理),本模块不再自持降级链
|
||
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "llm" / "scripts"))
|
||
from llm import chat_governed, extract_json # noqa: E402
|
||
from parse_outline import ensure_outline_coverage # noqa: E402
|
||
|
||
DSN = ("postgresql://root:f6710e2d0294eb1c10e26a805a64bc54@100.64.0.8:5433/muse-example"
|
||
"?keepalives=1&keepalives_idle=15&keepalives_interval=5&keepalives_count=3")
|
||
TENANT = 1
|
||
HERE = pathlib.Path(__file__).resolve().parent
|
||
TMP = pathlib.Path("/tmp/muse-parse")
|
||
OUTLINE_TARGET_RATIO = 0.045
|
||
MAX_OUTLINE_REPAIR_ATTEMPTS = 2
|
||
|
||
# ── 提示词资产(合同一律 load_contracts 动态渲染,禁手写——CONTRACTS 漂移冤案教训) ──
|
||
|
||
IDENTITY = """你是知识抽取员(extractor),分析槽位的默认绑定件。产出全部是草稿。
|
||
元数据纪律:schema 有什么字段你就抽什么,schema 没有的不抽——字段合同就是抽取 checklist,不自造结构;归型走各型「判据」;归不进任何型的候选=枚举缺口,如实报不硬塞;每字段要有正文证据,置信度低标「?」。
|
||
通则:以正文为准,不脑补正文没写的;基础字段规范填;采纳正文≠确认知识。"""
|
||
|
||
OUTLINE_COMPRESS_IDENTITY = """你是细纲压缩员。只压缩给定细纲,不补写剧情,不重新抽取其他内容。
|
||
必须严格遵守非空白字符上限,只输出指定 JSON 对象。"""
|
||
|
||
ENTITY_CRITERIA = ("实体型判据(脚手架级,只要 型/名称/一句话摘要):character=具名可指认的行动主体;"
|
||
"location=有名字的地点/星球/设施;faction=组织/国家/军团/公司;"
|
||
"power_system=力量体系/科技体系/修炼阶梯(体系本身,非招式);"
|
||
"item=有名字且有跨章戏份的装备/机甲/物品;event=已发生的重大事件(战役/事故/仪式)。"
|
||
"立卡门槛:有跨章戏份潜力;一次性龙套与单场景道具不收。")
|
||
|
||
PATTERN_TYPES = ("craft", "combat", "emotion", "scene_pattern", "trope")
|
||
|
||
# 基础公共字段(yudao 惯例列,出卡时由 ingest/payload 承载,不进抽取字段表)
|
||
BASE_KEYS = {"名称", "别名", "一句话摘要", "标签", "来源", "状态", "例证出处"}
|
||
|
||
# 抽象指代白名单(金标准结论:M3 自造代号如"A角色"破坏可读性;专名只允许进实例定位)
|
||
PLACEHOLDERS = "主角/对手/强敌/导师/盟友/队友/配角/反派/长辈/宝物/装备/机关/势力/秘密/危机"
|
||
|
||
|
||
def load_contracts(conn):
|
||
"""从库内 schema 版本快照动态渲染五型合同——prompt 与 ingest 守卫同源。
|
||
|
||
教训(B4 fable 审查坐实的 bug 级冤案):此前 prompt 手写合同与库内字段名漂移
|
||
(emotion/scene_pattern/trope 三型全对不上),M3 正确产出的 11 张非 craft 卡
|
||
全部被守卫按另一套字段名冤杀——「prompt 自立合同」违反元数据驱动公理,禁止再犯。
|
||
"""
|
||
contracts = {}
|
||
for t in PATTERN_TYPES:
|
||
# 走 active_version_id(激活版本=治理权威),与 ingest 守卫同版本语义——
|
||
# 「最新版本行」在铺了新版未激活时会与守卫劈叉
|
||
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]
|
||
fields = [f for f in snap.get("特有字段", []) if f.get("key") not in BASE_KEYS]
|
||
contracts[t] = {"中文名": snap.get("中文名", t), "判据": snap.get("判据", ""),
|
||
"字段": fields}
|
||
return contracts
|
||
|
||
|
||
def render_contracts(contracts):
|
||
"""合同渲染为 markdown 表格形态(比单行 JSON dump 对 LLM 遵从率更友好)。"""
|
||
parts = []
|
||
for t, c in contracts.items():
|
||
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)
|
||
|
||
|
||
def scaffold_prompt(title, ch, ch_title, text, prev_entities):
|
||
# 前文实体压缩(创始人 2026-07-14):判重只需名字清单,不必重发型+摘要全档——
|
||
# 全量 json 到后期每章膨胀到 15 万+ token(深空第 557 章 2851 实体≈15 万),是次数被压的真凶
|
||
_names = sorted({(e.get("名称") or e.get("name") or "").strip() for e in prev_entities} - {""}) if prev_entities else []
|
||
ents = "、".join(_names) if _names else "(第一章,为空)"
|
||
wc = len(text.replace("\n", "").replace(" ", ""))
|
||
cap = max(60, int(wc * 0.05)) # 绝对字数上限(纯比例约束 M3 会写超,首轮实测 9.8%–19%)
|
||
# 缓存友好(read-context「稳定度递减吃前缀缓存」):固定规则段全部前置吃前缀缓存,
|
||
# 章号/字数/上限/前文实体/正文(每章必变)一律尾置;cap 值随材料尾置、规则段只说「见文末上限」以保持固定
|
||
return f"""【功能指令(parse-book 章级 pass:细纲+实体+范式候选线索)】
|
||
对参考书的当前章(书名/章号/正文见文末材料区)做三件事:
|
||
1) 逆推本章细纲:章目标/关键事件/出场角色/伏笔动作(埋·推·收)/章末钩子。**细纲全文不得超过文末给出的字数上限**(约本章正文的 5%;超标会被校验脚本机械退回)。细纲是结构骨架,不是缩写复述——用短语与分号,不写完整句子。
|
||
2) 抽实体增量(脚手架级索引):{ENTITY_CRITERIA}
|
||
判重:文末给出前文已收录实体的**名字清单**,清单内的实体本章不重报(若本章赋予其新身份,可在摘要里点明);清单外的本章新实体照抽。
|
||
3) 报范式候选线索(**只报线索不出卡**,出卡由窗级聚类另做):本章表现突出、疑似可跨书复用的写法。五型候选:craft=单点叙事装置(删去它场景仍成立);combat=整场武力对抗的打法;emotion=整场情绪戏的推进;scene_pattern=拍卖/谈判/审讯等场景公式;trope=跨章复用的情节公式。
|
||
每条线索给四项:type(五型之一,拿不准填"?")、name(2–8 字短名,作者口头会说的话,禁书内专名禁修饰堆叠)、clue(一句话:这个写法怎么运作,用{PLACEHOLDERS}等抽象指代)、evidence(一句话:本章哪里这样写、为何突出,可用专名)。
|
||
0–5 条/章,平庸章 0 条正常;门槛比出卡低——拿不准的报上来,窗级聚类会过滤。
|
||
|
||
【输出规则(只输出一个 JSON 对象,禁止任何其他文字)】
|
||
{{"outline": "细纲文本", "entities": [{{"type": "六型之一", "name": "名称", "brief": "一句话摘要"}}], "hints": [{{"type": "五型之一或?", "name": "短名", "clue": "运作机制一句话", "evidence": "本章证据一句话"}}]}}
|
||
|
||
━━━ 本章材料(每章不同,非规则)━━━
|
||
【当前章】《{title}》第 {ch} 章《{ch_title}》|正文约 {wc} 字|细纲上限 {cap} 字
|
||
【前文实体名字清单(判重用,已在列的不重报)】
|
||
{ents}
|
||
【本章正文】
|
||
{text}"""
|
||
|
||
|
||
def window_cards_prompt(title, a, b, stage_outline, ch_outlines, ch_hints, contracts, batch_note):
|
||
"""窗级聚类出卡 prompt(B4-S3 核心资产;金标准七条纪律全落于此)。
|
||
放量教训:百条线索一次全型聚类=认知超载,M3 退化为逐条转写(37 章窗出 74 张卡)——
|
||
调用侧按型分批(batch_note 说明本批范围),每批线索几十条以内聚类才真实发生。
|
||
缓存友好(read-context「稳定度递减吃前缀缓存」):合同+七条纪律+输出格式(上千字、
|
||
跨书跨窗完全一致)全部前置吃前缀缓存;书名/章范围/本批说明/三层材料(每窗每批必变)尾置。"""
|
||
return f"""【功能指令(parse-book 窗级聚类出卡+脱敏红线)】
|
||
把一个窗口内已完成章级分析的候选线索**聚类归并**成「公共范式卡」——同一手法在多章的多次出现归并为一张母卡,带全部实例章号;线索不足处可从细纲的伏笔账(埋·推·收)补证跨章装置。书名、章范围、本批范围与三层材料(阶段大纲/逐章细纲/候选线索)均见文末材料区。
|
||
|
||
{render_contracts(contracts)}
|
||
|
||
出卡纪律:
|
||
- 聚类优先:先把线索按「同一运作机制」分组(不同章的同类线索=同一卡的多个实例),再逐组判断是否值得出卡。**本批最多 6 张,硬上限**——只出聚类后最能跨书复用的母卡;输出张数接近线索条数=你在逐条转写而不是聚类,整批作废。聚不成组又不够突出的线索直接丢弃。
|
||
- **每张卡必须带 instances**(窗内章号+一句话定位);缺 instances 的卡=机械拒收。
|
||
- 跨章证据:结构装置(伏笔类)尽量给出埋与收两端的实例章号;只有单章实例的会被标「跨窗待证」降低可信度。
|
||
- 命名:卡名=作者口头会说的话(2–8 字,如「先抑后扬」「借刀杀人」),禁「XX式YY化ZZ」修饰语堆叠、禁书内专名、禁自造黑话。
|
||
- 摘要:一句话**通用手法陈述**(抹掉本书信息后换任何题材仍成立),禁复述本书剧情、禁与实例定位雷同。
|
||
- 抽象指代白名单:{PLACEHOLDERS}——白名单外的自造代号(如"A角色""X道具")不要用。
|
||
- 可迁移性是卡的本体:说不清「作者这样做、读者会怎样」的因果=不是范式,整张卡不出;带此类字段(如原理)的型按其字段说明的句式写——**句式里的 X/Y/Z 是占位符,必须替换成本卡的具体内容**,字面保留「做Y」「读者会Z」=机械拒卡。
|
||
- 实例只报章号+一句话定位;**严禁在任何字段里自报间隔章数**——间隔由校验脚本按实例章号差机械计算,自己编数字=伪精确。
|
||
- fields 只允许使用该卡 type 在上方合同表里列出的中文 key,没证据的 key 省略不编造。**先按判据定 type,再对照该型的表格填字段**——把别的型才有的 key(哪怕内容再有价值)塞进本型=机械拒卡;上一批实测拒卡主因就是给非 craft 卡塞了 craft 才有的字段。
|
||
- 脱敏红线:只写抽象结构与手法归纳,严禁抄录原文(≥15 连续字与原文重合=机械拒卡);书内专名只允许出现在实例定位(anchor)里。
|
||
|
||
【输出规则(只输出一个 JSON 对象,禁止任何其他文字;没有值得立的卡时 cards 给空数组)】
|
||
{{"cards": [{{"type": "craft|combat|emotion|scene_pattern|trope", "name": "短名", "brief": "一句话通用摘要", "fields": {{"中文合同key": "值"}}, "instances": [{{"ch": 章号数字, "anchor": "一句话情节定位(可用专名)"}}]}}]}}
|
||
|
||
━━━ 本窗材料(每窗每批不同,非规则)━━━
|
||
《{title}》第 {a}–{b} 章。{batch_note}
|
||
|
||
【阶段大纲】
|
||
{stage_outline}
|
||
|
||
【逐章细纲】
|
||
{ch_outlines}
|
||
|
||
【逐章候选线索】
|
||
{ch_hints}"""
|
||
|
||
|
||
# 敏感降级链(历史常量):已并入 llm.BUDGET_CHAIN 全局统一链,由 chat_governed 内部消化,
|
||
# m3_json 不再自己遍历它;保留定义仅作历史留痕,勿再用于逐链换模型。
|
||
FALLBACK_MODELS = ["MiniMax-M2.7", "deepseek-v4-flash"]
|
||
|
||
|
||
class SensitiveHardStop(Exception):
|
||
"""主模型 + 两个降级备选全部撞敏感——按创始人指令硬停整个解析并汇报(不跳过、不硬扛)。"""
|
||
|
||
|
||
def m3_json(prompt, model, need_keys, system=IDENTITY):
|
||
"""调模型 → JSON 容错提取 → 形状校验(不合格带提示重试 1 次)。
|
||
|
||
敏感/额度降级已上收 llm.chat_governed(全局统一降级链 + 额度治理):撞内容安全或模型不可用
|
||
时由 chat_governed 沿全局链自动换模型、并按窗口预算/调用数治理;全链耗尽(used 为 None)时本
|
||
函数抛 SensitiveHardStop 交上层硬停(parse_upgrade/parse_salvage 靠它跳窗)。JSON 形状问题
|
||
非敏感,不换模型(原样抛 RuntimeError,上层按章跳过)。返回 (data, usage) 契约不变。"""
|
||
content, usage, used = chat_governed(prompt, model=model, system=system)
|
||
if used is None: # 全局降级链全部耗尽(内容安全/不可用)——保留硬停语义
|
||
raise SensitiveHardStop(f"chat_governed 全局降级链全部耗尽(内容安全或模型不可用),model={model}")
|
||
# 拿到内容:JSON 形状校验,不合格带提示重试 1 次(仍走 chat_governed 全局治理)
|
||
for retry in range(2):
|
||
try:
|
||
data = extract_json(content)
|
||
if all(k in data for k in need_keys):
|
||
return data, usage
|
||
err = f"缺少必需键 {need_keys}"
|
||
except Exception as e: # json_repair 也救不回来的输出
|
||
err = str(e)[:200]
|
||
if retry == 0:
|
||
content, u2, used2 = chat_governed(
|
||
prompt + f"\n\n【重试提示】上次输出无法解析({err}),请严格按输出规则只输出一个 JSON 对象。",
|
||
model=model, system=system)
|
||
if used2 is None:
|
||
break # 重试提示也触发全链耗尽:当作该轮失败,落到下方 RuntimeError
|
||
# 只合并数值键:上游 usage 带嵌套 dict(如 *_tokens_details),dict+dict 会崩(放量实测)
|
||
usage = {k: usage.get(k, 0) + u2.get(k, 0) for k in set(usage) | set(u2)
|
||
if isinstance(usage.get(k, 0), (int, float)) and isinstance(u2.get(k, 0), (int, float))}
|
||
# JSON 形状最终失败(非敏感):不换模型,原样抛(上层按章跳过)
|
||
raise RuntimeError(f"JSON 形状重试仍失败({used}): {err}")
|
||
|
||
|
||
def non_whitespace_len(text):
|
||
"""统计机械比例门使用的非空白字符数,保证调用前预检与 ingest 判据同口径。"""
|
||
return len(re.sub(r"\s", "", str(text or "")))
|
||
|
||
|
||
def outline_target_cap(source_text):
|
||
"""给细纲压缩留出相对 8% 硬门的安全余量,目标固定在正文的 4.5%。
|
||
|
||
60 字下限与 parse_ingest 的短章绝对豁免一致,避免感言、公告等极短章无解。
|
||
"""
|
||
return max(60, int(non_whitespace_len(source_text) * OUTLINE_TARGET_RATIO))
|
||
|
||
|
||
def outline_hard_cap(source_text):
|
||
"""返回与 parse_ingest 完全一致的 8% 最终硬门。"""
|
||
return max(60, int(non_whitespace_len(source_text) * 0.08))
|
||
|
||
|
||
def outline_compression_prompt(outline, cap, attempt):
|
||
"""构造只含细纲的短提示,不重发正文、实体名录或完整章级抽取任务。"""
|
||
current_length = non_whitespace_len(outline)
|
||
# 总上限被 M3 系统性忽略时,用逐短语预算再留约三成余量;仍只做语义压缩,不截字符串。
|
||
phrase_count = 4 if attempt == 1 else 2
|
||
phrase_cap = max(6, int(cap * (0.16 if attempt == 1 else 0.20)))
|
||
return f"""【细纲专用压缩|第 {attempt}/{MAX_OUTLINE_REPAIR_ATTEMPTS} 轮】
|
||
将下方细纲压成结构骨架,只保留章目标、关键事件、伏笔动作(埋/推/收)和章末钩子。
|
||
用短语与分号,删除修饰、对白、过程复述;不得新增原细纲没有的事实。
|
||
压缩结果的非空白字符不得超过 {cap},不得用空格或换行规避计数。
|
||
上一版共 {current_length} 个非空白字符。输出最多 {phrase_count} 个无标签短语,用分号连接;
|
||
每个短语不超过 {phrase_cap} 个非空白字符,总计仍不得超过 {cap}。不要写“目标:”“事件:”等标签。
|
||
不要重新执行其他抽取任务。只输出一个 JSON 对象,禁止任何其他文字:
|
||
{{"outline": "压缩后的细纲"}}
|
||
|
||
【待压缩细纲】
|
||
{outline}"""
|
||
|
||
|
||
def _merge_usage(total, current):
|
||
"""累加 token 数值项;忽略上游 usage 中不可相加的嵌套明细。"""
|
||
for key, value in (current or {}).items():
|
||
if isinstance(value, (int, float)):
|
||
total[key] = total.get(key, 0) + value
|
||
|
||
|
||
def repair_outline(first_data, source_text, model):
|
||
"""有限次数压缩首轮细纲;成功时仅替换 outline,其他首轮字段原样保留。
|
||
|
||
返回 ``(合并数据或 None, usage, 错误或 None)``。超长结果绝不机械截断,也不会
|
||
交给 ingest;调用方因此保留首轮比例门已写下的 failed 状态。
|
||
"""
|
||
cap = outline_target_cap(source_text)
|
||
hard_cap = outline_hard_cap(source_text)
|
||
candidate = str(first_data.get("outline") or "").strip()
|
||
usage_total = {}
|
||
last_length = non_whitespace_len(candidate)
|
||
last_error = None
|
||
for attempt in range(1, MAX_OUTLINE_REPAIR_ATTEMPTS + 1):
|
||
try:
|
||
compressed, usage = m3_json(
|
||
outline_compression_prompt(candidate, cap, attempt), model, ("outline",),
|
||
system=OUTLINE_COMPRESS_IDENTITY)
|
||
_merge_usage(usage_total, usage)
|
||
candidate = str(compressed.get("outline") or "").strip()
|
||
last_length = non_whitespace_len(candidate)
|
||
if candidate and (
|
||
last_length <= cap
|
||
or (attempt == MAX_OUTLINE_REPAIR_ATTEMPTS and last_length <= hard_cap)
|
||
):
|
||
repaired = dict(first_data)
|
||
repaired["outline"] = candidate
|
||
return repaired, usage_total, None
|
||
last_error = "输出为空" if not candidate else f"压缩输出 {last_length} 字,目标不超过 {cap} 字"
|
||
except RuntimeError as exc:
|
||
# JSON/调用错误允许进入下一次有限重试;敏感全链耗尽仍由 SensitiveHardStop 向上硬停。
|
||
last_error = f"调用失败:{exc}"
|
||
detail = last_error or f"压缩输出 {last_length} 字,目标不超过 {cap} 字"
|
||
return None, usage_total, (f"细纲专用压缩 {MAX_OUTLINE_REPAIR_ATTEMPTS} 轮仍失败:{detail};"
|
||
"未截断、未再次入库,保持 failed")
|
||
|
||
|
||
def ingest(kind, work_id, key, payload, keyflag="--chapter-order"):
|
||
"""写临时文件 → parse_ingest 机械校验入库;返回 (是否成功, 输出文本)。"""
|
||
TMP.mkdir(parents=True, exist_ok=True)
|
||
f = TMP / f"{work_id}-{key}-{kind}.json"
|
||
f.write_text(json.dumps(payload, ensure_ascii=False, indent=1))
|
||
r = subprocess.run([sys.executable, str(HERE / "parse_ingest.py"), kind,
|
||
"--work-id", str(work_id), keyflag, str(key), "--file", str(f)],
|
||
capture_output=True, text=True)
|
||
return r.returncode == 0, (r.stdout + r.stderr).strip()
|
||
|
||
|
||
def ingest_scaffold_with_repair(work_id, chapter_order, first_data, source_text, model):
|
||
"""首轮入库比例失败时只修细纲;返回入库结果、可追踪输出与压缩调用 usage。"""
|
||
ok, output = ingest("scaffold", work_id, chapter_order, first_data)
|
||
if ok or "细纲比例" not in output:
|
||
return ok, output, {}
|
||
repaired, usage, error = repair_outline(first_data, source_text, model)
|
||
if repaired is None:
|
||
return False, f"{output}\n{error}", usage
|
||
ok, repaired_output = ingest("scaffold", work_id, chapter_order, repaired)
|
||
return ok, repaired_output, usage
|
||
|
||
|
||
def _is_complete_scaffold_payload(data):
|
||
"""缓存只接受完整章级对象;数组成员也必须是对象,拒绝容错修复后的模糊形状。"""
|
||
return (
|
||
isinstance(data, dict)
|
||
and isinstance(data.get("outline"), str)
|
||
and bool(data["outline"].strip())
|
||
and isinstance(data.get("entities"), list)
|
||
and all(isinstance(item, dict) for item in data["entities"])
|
||
and isinstance(data.get("hints"), list)
|
||
and all(isinstance(item, dict) for item in data["hints"])
|
||
)
|
||
|
||
|
||
def resolve_scaffold_payload(work_id, chapter_order, scaffold_status, full_extract, trace_output=None):
|
||
"""仅为 failed 章复用严格缓存;其他情况惰性调用正文完整抽取。
|
||
|
||
返回 ``(载荷, usage, cache trace)``。缓存读取只用标准 JSON 解码,不走 LLM 输出的
|
||
容错提取;因此非法或残缺文件不会被误当成可复用首轮载荷。
|
||
"""
|
||
cache_path = TMP / f"{work_id}-{chapter_order}-scaffold.json"
|
||
|
||
def traced(message):
|
||
"""先落追踪输出再做可能耗时或失败的完整抽取。"""
|
||
if trace_output is not None:
|
||
trace_output(message)
|
||
return message
|
||
|
||
if scaffold_status != "failed":
|
||
trace = traced(f"cache miss: status={scaffold_status},不复用 {cache_path}")
|
||
data, usage = full_extract()
|
||
return data, usage, trace
|
||
try:
|
||
data = json.loads(cache_path.read_text(encoding="utf-8"))
|
||
except FileNotFoundError:
|
||
trace = traced(f"cache miss: 文件不存在 {cache_path}")
|
||
data, usage = full_extract()
|
||
return data, usage, trace
|
||
except (OSError, UnicodeError, json.JSONDecodeError) as exc:
|
||
trace = traced(f"cache miss: 无法严格解析 {cache_path}({type(exc).__name__})")
|
||
data, usage = full_extract()
|
||
return data, usage, trace
|
||
if not _is_complete_scaffold_payload(data):
|
||
trace = traced(f"cache miss: 对象或字段类型非法 {cache_path}")
|
||
fresh, usage = full_extract()
|
||
return fresh, usage, trace
|
||
trace = traced(f"cache hit: 复用首轮 failed 载荷 {cache_path}")
|
||
return data, {}, trace
|
||
|
||
|
||
def chapter_orders(conn, work_id, from_, to, incomplete_only):
|
||
"""返回本轮目标章;增量模式在循环前一次筛出 pending/failed,避免扫描全部 done 章。"""
|
||
if not incomplete_only:
|
||
return list(range(from_, to + 1))
|
||
rows = conn.execute(
|
||
"""SELECT c.order_no
|
||
FROM muse_content_chapter c
|
||
JOIN example_parse_task t ON t.chapter_id=c.id AND t.tenant_id=c.tenant_id
|
||
WHERE c.tenant_id=%s AND c.work_id=%s AND c.order_no BETWEEN %s AND %s
|
||
AND c.deleted=FALSE
|
||
AND COALESCE(t.scaffold_status, 'pending') IN ('pending', 'failed')
|
||
ORDER BY c.order_no""",
|
||
(TENANT, work_id, from_, to)).fetchall()
|
||
return [row[0] for row in rows]
|
||
|
||
|
||
def initialize_tasks(work_id, from_, to, incomplete_only):
|
||
"""首次全量运行建任务行;增量恢复复用已有任务表,避免再次逐章初始化。"""
|
||
if incomplete_only:
|
||
return
|
||
subprocess.run([sys.executable, str(HERE / "parse_ingest.py"), "init-tasks",
|
||
"--work-id", str(work_id), "--from", str(from_), "--to", str(to)],
|
||
capture_output=True, text=True)
|
||
|
||
|
||
@click.group()
|
||
def cli():
|
||
"""M3 直调拆书(章级线索 → 窗级出卡)"""
|
||
|
||
|
||
@cli.command()
|
||
@click.option("--work-id", type=int, required=True)
|
||
@click.option("--from", "from_", type=int, required=True)
|
||
@click.option("--to", type=int, required=True)
|
||
@click.option("--model", default="MiniMax-M3", show_default=True)
|
||
@click.option("--incomplete-only", is_flag=True,
|
||
help="只处理 pending/failed 章;循环前一次筛选,不逐章连接或打印 done 章")
|
||
def chapters(work_id, from_, to, model, incomplete_only):
|
||
"""章级 pass:逐章一次 M3(细纲+实体+范式候选线索)。正文只过这一遍。"""
|
||
initialize_tasks(work_id, from_, to, incomplete_only)
|
||
with psycopg.connect(DSN) as conn:
|
||
title = conn.execute("SELECT title FROM muse_content_work WHERE id=%s", (work_id,)).fetchone()[0]
|
||
targets = chapter_orders(conn, work_id, from_, to, incomplete_only)
|
||
if incomplete_only:
|
||
click.echo(f"incomplete-only: 待处理 {len(targets)} 章")
|
||
total_in = total_out = 0
|
||
cache_hits = cache_misses = 0
|
||
for ch in targets:
|
||
with psycopg.connect(DSN) as conn:
|
||
row = conn.execute(
|
||
"""SELECT c.id, c.title, b.content_text, t.scaffold_status
|
||
FROM muse_content_chapter c
|
||
JOIN muse_content_block b ON b.chapter_id=c.id AND b.deleted=FALSE
|
||
JOIN example_parse_task t ON t.chapter_id=c.id AND t.tenant_id=c.tenant_id
|
||
WHERE c.tenant_id=%s AND c.work_id=%s AND c.order_no=%s AND c.deleted=FALSE""",
|
||
(TENANT, work_id, ch)).fetchone()
|
||
if not row:
|
||
click.echo(f"#{ch} 章或任务行不存在,跳过")
|
||
continue
|
||
ch_id, ch_title, text, s_st = row
|
||
if s_st == "done":
|
||
click.echo(f"#{ch} 脚手架已完成,跳过")
|
||
continue
|
||
try:
|
||
def full_extract():
|
||
"""仅在缓存未命中时查询判重索引并执行正文完整抽取。"""
|
||
with psycopg.connect(DSN) as prev_conn:
|
||
prev = [e for (ents,) in prev_conn.execute(
|
||
"""SELECT s.entities FROM example_parse_scaffold s
|
||
JOIN muse_content_chapter c ON c.id=s.chapter_id
|
||
WHERE s.tenant_id=%s AND s.work_id=%s AND c.order_no<%s AND s.deleted=FALSE""",
|
||
(TENANT, work_id, ch)).fetchall() for e in ents]
|
||
prompt = scaffold_prompt(title, ch, ch_title, text, prev)
|
||
return m3_json(prompt, model, ("outline", "entities", "hints"))
|
||
|
||
if incomplete_only:
|
||
data, usage, cache_trace = resolve_scaffold_payload(
|
||
work_id, ch, s_st, full_extract,
|
||
trace_output=lambda trace: click.echo(f"#{ch} scaffold {trace}"))
|
||
if cache_trace.startswith("cache hit:"):
|
||
cache_hits += 1
|
||
else:
|
||
cache_misses += 1
|
||
else:
|
||
data, usage = full_extract()
|
||
total_in += usage.get("prompt_tokens", 0)
|
||
total_out += usage.get("completion_tokens", 0)
|
||
ok, out, repair_usage = ingest_scaffold_with_repair(work_id, ch, data, text, model)
|
||
total_in += repair_usage.get("prompt_tokens", 0)
|
||
total_out += repair_usage.get("completion_tokens", 0)
|
||
click.echo(f" {out}")
|
||
except SensitiveHardStop as e:
|
||
# 降级链(主+2 备)全撞敏感——按创始人指令硬停该书解析并汇报,不跳过、不硬扛
|
||
with psycopg.connect(DSN) as conn:
|
||
conn.execute(
|
||
"""UPDATE example_parse_task SET scaffold_status='failed', error_message=%s
|
||
WHERE tenant_id=%s AND work_id=%s AND chapter_id=(
|
||
SELECT id FROM muse_content_chapter
|
||
WHERE tenant_id=%s AND work_id=%s AND order_no=%s)""",
|
||
(str(e)[:500], TENANT, work_id, TENANT, work_id, ch))
|
||
conn.commit()
|
||
click.echo(f" ⛔ #{ch} 敏感降级链全失败:{e}")
|
||
click.echo(f"《{title}》解析按指令硬停于 #{ch}(此前完成章已落库;该章需人工定夺)")
|
||
return
|
||
except RuntimeError as e:
|
||
click.echo(f" #{ch} 章级 M3 失败: {e}")
|
||
if incomplete_only:
|
||
click.echo(f"scaffold cache: hit={cache_hits} miss={cache_misses}")
|
||
click.echo(f"《{title}》{from_}–{to} 章级完成;token in={total_in:,} out={total_out:,}")
|
||
|
||
|
||
@cli.command()
|
||
@click.option("--work-id", type=int, required=True)
|
||
@click.option("--from-order", "from_order", type=int, help="只出该窗(窗起始章);不给则全部窗")
|
||
@click.option("--model", default="MiniMax-M3", show_default=True)
|
||
@click.option("--redo", is_flag=True, help="窗内已有活卡也重出(默认跳过=断点续跑)")
|
||
def cards(work_id, from_order, model, redo):
|
||
"""窗级出卡:逐窗一次 M3 聚类归并(窗=example_parse_outline 行,先跑 parse_outline window)。"""
|
||
with psycopg.connect(DSN) as conn:
|
||
title = conn.execute("SELECT title FROM muse_content_work WHERE id=%s", (work_id,)).fetchone()[0]
|
||
contracts = load_contracts(conn) # prompt 与 ingest 守卫同源(库内 schema 快照)
|
||
sql = """SELECT from_order, to_order, outline_text, window_no FROM example_parse_outline
|
||
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE"""
|
||
args = [TENANT, work_id]
|
||
if from_order:
|
||
sql += " AND from_order=%s"
|
||
args.append(from_order)
|
||
wins = conn.execute(sql + " ORDER BY from_order", args).fetchall()
|
||
if not wins:
|
||
raise click.ClickException("无大纲窗行——先跑 parse_outline.py window(窗是出卡的切分依据)")
|
||
if from_order is None:
|
||
with psycopg.connect(DSN) as conn:
|
||
bounds = conn.execute(
|
||
"""SELECT min(order_no), max(order_no) FROM muse_content_chapter
|
||
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""",
|
||
(TENANT, work_id),
|
||
).fetchone()
|
||
if bounds and bounds[0] is not None:
|
||
ensure_outline_coverage(
|
||
[(row[0], row[1]) for row in wins],
|
||
bounds[0],
|
||
bounds[1],
|
||
context="cards 前",
|
||
)
|
||
total_in = total_out = 0
|
||
for a, b, stage_ol, wno in wins:
|
||
with psycopg.connect(DSN) as conn:
|
||
# 断点续跑:窗内已有活卡(本窗出的)则跳过
|
||
if not redo and conn.execute(
|
||
"""SELECT 1 FROM muse_knowledge_draft WHERE tenant_id=%s AND source_type='parse_book'
|
||
AND source_id=%s AND draft_payload->>'窗起'=%s AND deleted=FALSE LIMIT 1""",
|
||
(TENANT, work_id, str(a))).fetchone():
|
||
click.echo(f"win#{a}–{b} 已有活卡,跳过(--redo 强制重出)")
|
||
continue
|
||
rows = conn.execute(
|
||
"""SELECT c.order_no, c.title, s.outline_text, s.pattern_hints
|
||
FROM muse_content_chapter c
|
||
JOIN example_parse_scaffold s ON s.chapter_id=c.id AND s.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()
|
||
if len(rows) < (b - a + 1):
|
||
have = {r[0] for r in rows}
|
||
click.echo(f"win#{a}–{b} 缺章级脚手架 {[n for n in range(a, b + 1) if n not in have][:10]},先补 chapters")
|
||
continue
|
||
ch_outlines = "\n".join(f"第{no}章《{ct}》:{o}" for no, ct, o, _ in rows)
|
||
hints = [{"章": no, "线索": h} for no, _, _, hs in rows for h in (hs or [])]
|
||
n_hint = len(hints)
|
||
# 按型分批出卡(放量实测:百条线索一次全型聚类=认知超载,退化为逐条转写)——
|
||
# 每批只给该型合同+该型线索,型内聚类真实可行;「?」型线索单独一批(五型合同全给,让 M3 先归型)
|
||
by_type = {}
|
||
for h in hints:
|
||
t = h["线索"].get("型")
|
||
by_type.setdefault(t if t in PATTERN_TYPES else "?", []).append(h)
|
||
all_cards, win_in, win_out = [], 0, 0
|
||
for t, hs in sorted(by_type.items()):
|
||
sub_contracts = contracts if t == "?" else {t: contracts[t]}
|
||
note = (f"本批只处理「{t}」型的 {len(hs)} 条线索(其余型另批处理),只出 {t} 型卡。"
|
||
if t != "?" else
|
||
f"本批是 {len(hs)} 条章级未归型的线索——先按各型判据归型,归不进五型的丢弃。")
|
||
try:
|
||
p = window_cards_prompt(title, a, b, stage_ol, ch_outlines,
|
||
json.dumps(hs, ensure_ascii=False, indent=1), sub_contracts, note)
|
||
data, usage = m3_json(p, model, ("cards",))
|
||
win_in += usage.get("prompt_tokens", 0)
|
||
win_out += usage.get("completion_tokens", 0)
|
||
# 病症机械检测→带方重试一次(放量实测三病:方差空批/逐条转写/长清单丢实例键)
|
||
cs = data.get("cards") or []
|
||
no_ins = sum(1 for c in cs if not (c.get("instances") or c.get("实例")))
|
||
sick = ("输出了空数组,但本批有 %d 条候选线索——认真逐条聚类评估后再输出,确实全部不够格才允许空数组"
|
||
% len(hs) if (not cs and len(hs) >= 3) else
|
||
"输出了 %d 张卡——这是逐条转写不是聚类。重新按「同一运作机制」分组归并,只出最能跨书复用的母卡,最多 6 张"
|
||
% len(cs) if len(cs) > 8 else
|
||
"有 %d 张卡缺 instances(窗内章号+定位)——instances 是必填键,缺=全部被机械拒收。重新输出全部卡并逐张补上"
|
||
% no_ins if no_ins else None)
|
||
if sick:
|
||
data, usage = m3_json(p + "\n\n【重试】上次" + sick + "。", model, ("cards",))
|
||
win_in += usage.get("prompt_tokens", 0)
|
||
win_out += usage.get("completion_tokens", 0)
|
||
cs = data.get("cards") or []
|
||
all_cards += cs
|
||
click.echo(f" 批[{t}] 线索{len(hs)} → 卡{len(cs)}")
|
||
except RuntimeError as e:
|
||
# m3_json 已完成调用/结构重试;继续会让该型线索永久丢失并把整窗误标 done。
|
||
raise click.ClickException(
|
||
f"win#{a}–{b} 批[{t}] 连续重试仍失败,停止当前作品并保留断点:{str(e)[:180]}"
|
||
) from e
|
||
total_in += win_in
|
||
total_out += win_out
|
||
_, out = ingest("cards", work_id, a, {"cards": all_cards}, keyflag="--from-order")
|
||
click.echo(f"win#{a}–{b}(线索 {n_hint} 条,分 {len(by_type)} 批)\n {out}")
|
||
click.echo(f"《{title}》窗级出卡完成;token in={total_in:,} out={total_out:,}")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
try:
|
||
cli()
|
||
except (psycopg.Error, RuntimeError) as e:
|
||
click.echo(f"[错误] {type(e).__name__}: {e}", err=True)
|
||
sys.exit(1)
|