zizi 9c87f0dee9 feat(parse-book): 范式卡回溯聚类去重(recluster)——同型近义卡嵌入初筛+M3终判归并母卡,全库7045→6327(并718/~10%)
recluster 子命令:连接三段(LLM期不持DB连接)/贪心最近母卡(避union-find巨簇)/软删可逆/dry-run门控;母卡累积实例+回溯归并审计(被并卡全文留档可回滚)。前向修复 cards():SIM_MERGE 0.85→0.80、最近邻≥0.75即送 merge_judge(堵住跨书近义卡从不判定的病根)。

五型放量:combat886→770/emotion1151→1027/craft1445→1326/scene_pattern1753→1578/trope1810→1626,并718。真实归并率~10%(向量投影74%多为同型词汇假近,merge_judge剪掉85-95%;创始人知情拍板全量0.75)。可逆性/完整性/软删/质量四项主代理亲验。

llm.py:睡窗日志边界显示 typo 修复(secs取整落在04:59:59被%H显示成非法边界04:00,+1s归整为05:00)。
2026-07-17 06:10:04 +08:00

736 lines
45 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 配套确定性脚本:拆书产物校验入库 + 任务状态机(B2/B4-S3)。
职责边界:extractor(LLM) 只产结构化 JSON 文件,不碰库;本脚本做机械校验后写库——
字段 key 合法性(对库内字段合同)、五型归型、出处必填、**脱敏红线 15 连字检测**(硬阻断)。
状态全在库(example_parse_task),断点续跑与幂等按章/按窗。
B4-S3 起管线两级:scaffold(章级:细纲+实体+范式候选线索)→ cards(窗级:聚类母卡)。
本脚本仅在两个受控点调用 LLM/嵌入服务(不产内容):
① 新卡嵌入(判重与检索基座共用);② 相似度 ≥0.85 时的归并判定(输入两卡 JSON,输出 merge/keep)。
patterns 命令是 S3 前的章级出卡入口,保留作回滚保险,新管线不再使用。
"""
import hashlib
import json
import pathlib
import re
import sys
import click
import psycopg
from psycopg.types.json import Jsonb
# 受控点依赖:嵌入走 embed skill 同一实现(同模型同维),归并判定走 llm skill 统一入口
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "embed" / "scripts"))
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "llm" / "scripts"))
from embed_drafts import _session, build_embed_text, embed_texts # noqa: E402
from llm import chat_governed, extract_json # 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, ACTOR = 1, "1"
PATTERN_TYPES = {"craft", "combat", "emotion", "scene_pattern", "trope"} # 拍板①:首轮只拆五型
NGRAM = 15 # 脱敏红线:≥15 连续字与原文重合=违规(parse-book skill)
# 版权 IP 与系统专名词表(复抽实证:机战无限为高达系同人,IP 背景词不在实体名池)。
# 词形边界(opus 终检教训):
# - "浮游炮"不入表——已泛化为通用武器品类词(同"光剑"),入表实测误伤 3 卡;
# - "战功"/"负能"不入表——通用词(立下战功/负能量)子串误伤面大,其系统义
# (机战货币"战功"、星环设定"负能")由批次清洗按语境甄别;
# - 原创书自造宇宙观术语(伪造物主/幽畸/锐眼)同属泄漏,一并列管。
IP_LEAK_WORDS = ("GN粒子", "太阳炉", "扎古", "高达", "夏亚", "阿姆罗", "米诺夫斯基",
"脑量子波", "GN-Bit", "影印人", "殖装", "战功点", "宇宙世纪",
"爆种", "次元兽", "伪造物主", "幽畸", "锐眼")
EMBED_MODEL = "Qwen/Qwen3-Embedding-8B"
SIM_MERGE = 0.80 # 强候选线:≥该值=高置信近义(前向修复后 0.75+ 均送 merge_judge,本值只作档位标注)
SIM_MARK = 0.75 # 候选线:≥该值即送 merge_judge 终判(cards 前向判重 + recluster 回溯聚类 同一判据)
def norm_scaffold(data: dict) -> dict:
"""键名归一:StructuredOutput 工具 schema 只允许 ASCII 键(API 硬约束 2026-07-13 实测),
库内 payload 统一中文键(供下游 prompt 消费)。两种键名都接受。"""
ents = []
for e in data.get("实体") or data.get("entities") or []:
ents.append({"型": e.get("型") or e.get("type"), "名称": e.get("名称") or e.get("name"),
"一句话摘要": e.get("一句话摘要") or e.get("brief")})
hints = []
for h in data.get("线索") or data.get("hints") or []:
hints.append({"型": h.get("型") or h.get("type"), "短名": h.get("短名") or h.get("name"),
"线索": h.get("线索") or h.get("clue"), "证据": h.get("证据") or h.get("evidence")})
return {"细纲": data.get("细纲") or data.get("outline") or "", "实体": ents, "线索": hints}
def norm_card(c: dict) -> dict:
"""范式卡键名归一(ASCII→中文);窗级卡带实例数组 [{章,定位}]。"""
src = c.get("出处") or c.get("source") or {}
inss = []
for ins in c.get("实例") or c.get("instances") or []:
try:
ch_no = int(ins.get("章") or ins.get("ch"))
except (TypeError, ValueError):
continue # 章号非数字=编造,丢弃该实例
inss.append({"章": ch_no, "定位": ins.get("定位") or ins.get("anchor") or ""})
return {"型": c.get("型") or c.get("type"), "名称": c.get("名称") or c.get("name"),
"一句话摘要": c.get("一句话摘要") or c.get("brief"),
"字段": c.get("字段") or c.get("fields") or {},
"实例": inss,
"出处": {"书名": src.get("书名") or src.get("book"),
"回目": src.get("回目") or src.get("chapter"),
"定位": src.get("定位") or src.get("anchor")}}
def chapter_of(conn, work_id, order_no):
"""取章 id 与正文。"""
row = conn.execute(
"""SELECT c.id, 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.order_no=%s AND c.deleted=FALSE""",
(TENANT, work_id, order_no)).fetchone()
if not row:
raise click.ClickException(f"章不存在: work={work_id} order={order_no}")
return row
def field_contract(conn, ttype):
"""库内字段合同:型 → 合法字段 key 集合。"""
rows = conn.execute(
"""SELECT f.field_key FROM muse_meta_field f
JOIN muse_meta_schema_version sv ON sv.id=f.schema_version_id
JOIN muse_meta_schema s ON s.active_version_id=sv.id
WHERE s.tenant_id=%s AND s.schema_key=%s AND f.deleted=FALSE""",
(TENANT, ttype)).fetchall()
return {r[0] for r in rows}
def load_entity_names(conn, work_id):
"""本书专名词典:scaffold 实体名就是现成专名表(短专名泄漏 15 连字抓不到,
金标准抓到 4 张「海拉/千亿夸克」级泄漏);出处/实例定位字段豁免(设计允许专名只进出处)。"""
names = {e.get("名称") or e.get("name") for (ents,) in conn.execute(
"""SELECT entities FROM example_parse_scaffold
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""",
(TENANT, work_id)).fetchall() for e in ents}
return {n for n in names if n and len(n) >= 2}
def leak_check(card_texts, source_text):
"""脱敏机械检查:卡内任一文本值含与原文 ≥NGRAM 连续字重合 → 返回违规片段。"""
src = re.sub(r'\s', '', source_text)
grams = {src[i:i + NGRAM] for i in range(0, max(0, len(src) - NGRAM + 1))}
for t in card_texts:
tt = re.sub(r'\s', '', str(t))
for i in range(0, max(0, len(tt) - NGRAM + 1)):
if tt[i:i + NGRAM] in grams:
return tt[i:i + NGRAM]
return None
def validate_card(c, contracts, entity_names, src_text):
"""卡片通用机械守卫(章级 patterns 与窗级 cards 共用——守卫同源,禁两处各写一套)。
返回 (拒卡原因列表, 改型审计|None);c 可能被就地改型。"""
t = c.get("型")
reasons, retyped = [], None
if t not in PATTERN_TYPES:
reasons.append(f"型不合法:{t}(首轮只拆五型)")
if not c.get("名称"):
reasons.append("缺名称")
fields = c.get("字段") or {}
if t in contracts and fields:
illegal = set(fields) - contracts[t]
if illegal:
# 字段指纹改型:LLM 的 type 标注不可靠(实测 M3 系统性全标 craft),
# 五型字段 key 互不重叠——字段集合全命中唯一他型合同时,以字段为准改型
fits = [t2 for t2 in PATTERN_TYPES
if t2 != t and fields and set(fields) <= contracts[t2]]
if len(fits) == 1:
retyped = f"{t}→{fits[0]}(字段指纹改型)"
c["型"] = t = fits[0]
else:
# 裁剪降级(B4-S3 实测:M3 往 trope 塞「读者收益」类他型增益字段屡教不改):
# 越界字段剥离(内容在本型无处安放,不该陪葬整卡),审计留痕;
# 裁掉必填字段的情况由下方 craft 双模板必填校验兜底拒卡
for k in illegal:
fields.pop(k)
c["裁剪字段"] = sorted(illegal)
# 守卫(B4-S2):craft 双模板条件必填/禁填 + 装置类型闭合枚举 + 禁复合标签
if t == "craft" and not reasons:
form = (fields.get("装置形态") or "").strip()
dtype = (fields.get("装置类型") or "").strip()
if form not in ("结构装置", "场景手法"):
reasons.append(f"装置形态非法:{form!r}(须为 结构装置/场景手法)")
elif form == "结构装置":
miss = [k for k in ("埋设手法", "回收点", "记忆维持", "间隔纪律") if not fields.get(k)]
if miss:
reasons.append(f"结构装置必填缺失:{miss}")
else: # 场景手法:禁填间隔纪律(同场景闭环没有间隔,硬填必出伪纪律)
miss = [k for k in ("复用节奏", "异质化要求", "单章上限") if not fields.get(k)]
if miss:
reasons.append(f"场景手法必填缺失:{miss}")
if fields.get("间隔纪律"):
reasons.append("场景手法禁填间隔纪律(伪纪律来源)")
if dtype:
if any(sep in dtype for sep in ("+", "、", "/", ",")):
reasons.append(f"装置类型禁复合标签:{dtype!r}(只填最主要的一个)")
elif dtype not in ("伏笔", "契诃夫之枪", "信息差", "重复意象", "倒计时", "身份错认", "非装置"):
reasons.append(f"装置类型不在闭合枚举:{dtype!r}")
# 原理句式占位符检测(实测病:M3 把模板「当X时做Y因为读者会Z」的占位符字面抄进卡)
if re.search(r"[做当会][XYZ]", str(fields.get("原理") or "")):
reasons.append("原理含未展开的模板占位符(X/Y/Z 须替换为具体内容)")
# 脱敏红线 + 专名词典(实例定位/出处豁免——设计允许专名只进定位)
texts = [x for x in (c.get("名称"), c.get("一句话摘要"), *fields.values()) if x]
leak = leak_check(texts, src_text)
if leak:
reasons.append(f"脱敏违规(≥{NGRAM}连字重合):「{leak}」")
hit_names = [n for n in entity_names if any(n in str(x) for x in texts)]
if hit_names:
reasons.append(f"专名泄漏(本书实体名):{hit_names[:3]}")
# 版权 IP 与系统专名词表(fable5 复抽实证盲区:同人书的 IP 背景词不在 scaffold
# 实体名里,批7 全部 12 张硬泄漏均由此逃逸,其中 4 张专名直入卡名——texts 含
# 卡名,此处一并覆盖。词表按实证泄漏词维护,勿加过泛词防误伤)
hit_ip = [w for w in IP_LEAK_WORDS if any(w in str(x) for x in texts)]
if hit_ip:
reasons.append(f"专名泄漏(版权IP/系统词):{hit_ip[:3]}")
return reasons, retyped
def set_task(conn, work_id, chapter_id, **cols):
"""推进任务状态机(attempt 自增)。"""
sets = ", ".join(f"{k}=%s" for k in cols)
conn.execute(
f"""UPDATE example_parse_task SET {sets}, attempt_count=attempt_count+1, updater=%s
WHERE tenant_id=%s AND work_id=%s AND chapter_id=%s""",
(*cols.values(), ACTOR, TENANT, work_id, chapter_id))
@click.group()
def cli():
"""拆书入库与任务状态机"""
@cli.command("init-tasks")
@click.option("--work-id", type=int, required=True)
@click.option("--from", "from_", type=int, default=1, show_default=True)
@click.option("--to", type=int, required=True)
def init_tasks(work_id, from_, to):
"""按章建任务行(幂等),并把参考书档案 parse_scope/parse_status 置为拆书中。"""
with psycopg.connect(DSN) as conn:
chs = conn.execute(
"""SELECT id, order_no FROM muse_content_chapter
WHERE tenant_id=%s AND work_id=%s AND order_no BETWEEN %s AND %s AND deleted=FALSE
ORDER BY order_no""", (TENANT, work_id, from_, to)).fetchall()
n = 0
for ch_id, _no in chs:
conn.execute(
"""INSERT INTO example_parse_task (work_id, chapter_id, creator, updater, tenant_id)
VALUES (%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, work_id, chapter_id) DO NOTHING""",
(work_id, ch_id, ACTOR, ACTOR, TENANT))
n += 1
conn.execute(
"""UPDATE example_reference_work SET parse_scope=%s, parse_status='parsing', updater=%s
WHERE tenant_id=%s AND work_id=%s""",
(Jsonb({"from": from_, "to": to}), ACTOR, TENANT, work_id))
conn.commit()
click.echo(f"任务行就绪: work={work_id} 章 {from_}–{to}({n} 行)")
@cli.command()
@click.option("--work-id", type=int, required=True)
@click.option("--chapter-order", type=int, required=True)
@click.option("--file", "file_", type=click.Path(exists=True), required=True)
def scaffold(work_id, chapter_order, file_):
"""章级入库:{细纲, 实体:[{型,名称,一句话摘要}], 线索:[{型,短名,线索,证据}]};比例约束校验。"""
data = norm_scaffold(json.loads(pathlib.Path(file_).read_text()))
with psycopg.connect(DSN) as conn:
ch_id, src = chapter_of(conn, work_id, chapter_order)
outline = (data.get("细纲") or "").strip()
if not outline:
raise click.ClickException("细纲为空")
ratio = len(re.sub(r'\s', '', outline)) / max(1, len(re.sub(r'\s', '', src)))
# 8% 防的是「摘要化伪装结构化」(长章复述);≤60 字的细纲对任何章都是高度压缩、
# 不可能是复述,绝对豁免——否则极短章(感言/公告,实测 116 字)永远无法过闸
if ratio > 0.08 and len(re.sub(r'\s', '', outline)) > 60: # 拍板值 3–5%,8% 为机械硬顶
set_task(conn, work_id, ch_id, scaffold_status="failed",
error_message=f"细纲比例超标 {ratio:.1%}>8%")
conn.commit()
raise click.ClickException(f"细纲比例 {ratio:.1%} 超标(>8%),已记失败退回重解析")
# 实体坏行宽容(放量实测:M3 偶发报 名称=null 的实体,一条坏行不该陪葬整章——
# 细纲/线索/其余实体全丢且任务无失败记录,窗切分会报缺章)
ents, n_bad_ent = [], 0
for e in data.get("实体") or []:
if e.get("名称") and e.get("型"):
ents.append(e)
else:
n_bad_ent += 1
# 范式候选线索(B4-S3):候选层宽容处置——坏行丢弃计数,不整批拒
hints, n_bad = [], 0
for h in data.get("线索") or []:
if h.get("短名") and h.get("线索") and (h.get("型") in PATTERN_TYPES or h.get("型") == "?"):
hints.append(h)
else:
n_bad += 1
conn.execute(
"""INSERT INTO example_parse_scaffold (work_id, chapter_id, outline_text, entities,
pattern_hints, creator, updater, tenant_id)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, chapter_id)
DO UPDATE SET outline_text=EXCLUDED.outline_text, entities=EXCLUDED.entities,
pattern_hints=EXCLUDED.pattern_hints, deleted=FALSE, updater=EXCLUDED.updater""",
(work_id, ch_id, outline, Jsonb(ents), Jsonb(hints), ACTOR, ACTOR, TENANT))
set_task(conn, work_id, ch_id, scaffold_status="done", error_message=None)
conn.commit()
click.echo(f"scaffold✓ work={work_id} ch#{chapter_order}: 细纲{len(outline)}字({ratio:.1%}) "
f"实体{len(ents)}" + (f"(弃坏行{n_bad_ent})" if n_bad_ent else "")
+ f" 线索{len(hints)}" + (f"(弃坏行{n_bad})" if n_bad else ""))
@cli.command()
@click.option("--work-id", type=int, required=True)
@click.option("--chapter-order", type=int, required=True)
@click.option("--file", "file_", type=click.Path(exists=True), required=True)
def patterns(work_id, chapter_order, file_):
"""【旧管线·回滚保险】章级范式卡入库;B4-S3 后由窗级 cards 替代,新管线勿用。"""
raw = json.loads(pathlib.Path(file_).read_text())
if isinstance(raw, dict): # 兼容 {cards:[…]}/{卡:[…]} 包装
raw = raw.get("cards") or raw.get("卡") or []
if not isinstance(raw, list):
raise click.ClickException("patterns 文件须为卡片数组")
cards = [norm_card(c) for c in raw]
with psycopg.connect(DSN) as conn:
ch_id, src = chapter_of(conn, work_id, chapter_order)
contracts = {t: field_contract(conn, t) for t in PATTERN_TYPES}
book = conn.execute("SELECT title FROM muse_content_work WHERE id=%s", (work_id,)).fetchone()[0]
entity_names = load_entity_names(conn, work_id)
# 幂等:重跑本章 = 软删本章旧 draft。
# 但新轮 0 卡时保留旧卡——LLM 判卡有随机性,重跑「0 卡」不应清掉上轮已验证的产出
if cards:
conn.execute(
"""UPDATE muse_knowledge_draft SET deleted=TRUE, updater=%s
WHERE tenant_id=%s AND source_type='parse_book' AND source_id=%s
AND draft_payload->>'章序'=%s AND deleted=FALSE""",
(ACTOR, TENANT, work_id, str(chapter_order)))
ok, rejected = 0, []
for i, c in enumerate(cards):
reasons, retyped = validate_card(c, contracts, entity_names, src)
src_ref = c.get("出处") or {}
if not (src_ref.get("书名") and src_ref.get("回目")):
reasons.append("出处不完整(需书名+回目)")
if reasons:
rejected.append({"卡": c.get("名称") or f"#{i}", "原因": reasons})
continue
payload = {"型": c["型"], "名称": c["名称"], "一句话摘要": c.get("一句话摘要", ""),
"字段": c.get("字段") or {}, "出处": src_ref, "目标库": "公共范式库",
"章序": chapter_order, "来源": f"拆书@{book}", "状态": "草稿"}
if retyped:
payload["改型"] = retyped # 审计:机械改型可追溯
if c.get("裁剪字段"):
payload["裁剪字段"] = c["裁剪字段"] # 审计:越合同字段被剥离入库
cid = f"parse-{work_id}-{chapter_order}-{i}-" + hashlib.sha256(
json.dumps(payload, ensure_ascii=False, sort_keys=True).encode()).hexdigest()[:8]
conn.execute(
"""INSERT INTO muse_knowledge_draft (work_id, draft_type, draft_payload, status,
source_type, source_id, command_id, creator, updater, tenant_id)
VALUES (0,'entity',%s,'pending','parse_book',%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, command_id) WHERE command_id IS NOT NULL DO NOTHING""",
(Jsonb(payload), work_id, cid, ACTOR, ACTOR, TENANT))
ok += 1
set_task(conn, work_id, ch_id, pattern_status="done",
error_message=None if not rejected else f"拒卡{len(rejected)}: " + json.dumps(rejected, ensure_ascii=False)[:900])
conn.commit()
click.echo(f"patterns✓ work={work_id} ch#{chapter_order}: 入库{ok} 拒{len(rejected)}")
for r in rejected:
click.echo(f" [拒] {r['卡']}: {'; '.join(r['原因'])}")
def merge_judge(new_payload, old_payload, model="MiniMax-M3"):
"""受控 LLM 点②:嵌入初筛 ≥SIM_MERGE 后的归并终判(输入两卡 JSON,输出 merge/keep)。
双保险设计:embedding 只做初筛,是否同一手法由 LLM 判——单靠阈值必产生错误合并。"""
prompt = ("两张「公共范式卡」由不同窗/不同书各自独立归纳(JSON 附后)。判断它们是否为同一写作手法:\n"
"- 运作机制相同(哪怕载体/题材/叫法不同)→ merge;\n"
"- 机制不同、或适用场景本质不同 → keep。拿不准时选 keep(错误合并比重复卡更伤)。\n"
'输出规则(只输出一个 JSON 对象):{"verdict": "merge"或"keep", "reason": "一句话依据"}\n\n'
f"【卡A(已入库)】\n{json.dumps(old_payload, ensure_ascii=False)}\n\n"
f"【卡B(新卡)】\n{json.dumps(new_payload, ensure_ascii=False)}")
try:
content, _, used = chat_governed(prompt, model=model)
# 全局额度链全部拦截(content 为 None):保守判「keep」=不合并——判重失败绝不能误并,
# 也不该因此卡住或崩书(错误合并比重复卡更伤,宁重复不误并)
if content is None:
return "keep", "归并判定全链耗尽(内容安全/不可用),保守保留不合并"
data = extract_json(content)
return data.get("verdict"), data.get("reason", "")
except Exception as e: # 判定失败=keep(不阻断入库,宁重复不误并)
return "keep", f"归并判定调用失败({type(e).__name__}),默认保留"
@cli.command()
@click.option("--work-id", type=int, required=True)
@click.option("--from-order", "from_order", type=int, required=True, help="窗起始章(对应大纲窗行)")
@click.option("--file", "file_", type=click.Path(exists=True), required=True)
@click.option("--model", default="MiniMax-M3", show_default=True, help="归并判定用模型")
def cards(work_id, from_order, file_, model):
"""窗级母卡入库(B4-S3):{cards:[{型,名称,摘要,字段,实例[{章,定位}]}]}。
守卫=章级全部 + 窗级新增:实例域校验 / 间隔章数机械计算(伪精确灭绝)/ 嵌入判重→归并判定。"""
raw = json.loads(pathlib.Path(file_).read_text())
if isinstance(raw, dict):
raw = raw.get("cards") or raw.get("卡") or []
cards_in = [norm_card(c) for c in raw]
with psycopg.connect(DSN) as conn:
win = conn.execute(
"""SELECT to_order FROM example_parse_outline
WHERE tenant_id=%s AND work_id=%s AND from_order=%s AND deleted=FALSE""",
(TENANT, work_id, from_order)).fetchone()
if not win:
raise click.ClickException(f"窗不存在: work={work_id} from={from_order}——先跑 parse_outline window")
to_order = win[0]
book = conn.execute("SELECT title FROM muse_content_work WHERE id=%s", (work_id,)).fetchone()[0]
# 窗内正文拼接(15 连字红线对照源:卡可能引到窗内任何一章)
src = "".join(t for (t,) in conn.execute(
"""SELECT 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.order_no BETWEEN %s AND %s
AND c.deleted=FALSE ORDER BY c.order_no""",
(TENANT, work_id, from_order, to_order)).fetchall())
contracts = {t: field_contract(conn, t) for t in PATTERN_TYPES}
entity_names = load_entity_names(conn, work_id)
# 先全量校验,收集通过集——软删旧卡延后到「确有新卡通过」才执行。
# 教训(放量实测):软删在校验前时,新卡全拒(如 M3 整批丢实例键)=旧卡已删+新卡全拒=净损失
passed, rejected = [], []
for i, c in enumerate(cards_in):
reasons, retyped = validate_card(c, contracts, entity_names, src)
# 实例域校验:章号必须落在本窗内——越窗章号=编造(跨窗同手法靠嵌入判重归并)
inss = [ins for ins in c.get("实例") or [] if from_order <= ins["章"] <= to_order]
n_drop = len(c.get("实例") or []) - len(inss)
if not inss:
reasons.append("实例为空或章号全部越窗(窗级卡必须带窗内实例)")
fields = c.get("字段") or {}
# 间隔章数机械计算(伪精确灭绝):结构装置实例≥2 → 章号差为权威值追加;
# 单实例结构装置=本窗证据不全(可能收在后续窗),标「跨窗待证」入库不拒
cross_pend = False
if not reasons and c.get("型") == "craft" and fields.get("装置形态") == "结构装置":
ch_nos = sorted({ins["章"] for ins in inss})
if len(ch_nos) >= 2:
fields["间隔纪律"] = ((fields.get("间隔纪律") or "").strip()
+ f"(实测:#{ch_nos[0]}埋→#{ch_nos[-1]}收,隔{ch_nos[-1] - ch_nos[0]}章)").strip()
else:
cross_pend = True
if reasons:
rejected.append({"卡": c.get("名称") or f"#{i}", "原因": reasons})
continue
payload = {"型": c["型"], "名称": c["名称"], "一句话摘要": c.get("一句话摘要", ""),
"字段": fields, "实例": inss,
"出处": {"书名": book, "回目": f"第{inss[0]['章']}章", "定位": inss[0]["定位"]},
"目标库": "公共范式库", "窗起": from_order, "来源": f"拆书@{book}", "状态": "草稿"}
if retyped:
payload["改型"] = retyped
if c.get("裁剪字段"):
payload["裁剪字段"] = c["裁剪字段"] # 审计:越合同字段被剥离入库
if cross_pend:
payload["跨窗待证"] = True # 审核环可见:结构装置只有单端实例
if n_drop:
payload["越窗实例弃"] = n_drop # 审计:M3 报了窗外章号
passed.append((i, payload, inss))
# 确有通过卡才软删本窗旧卡(幂等重出;全拒/空批不动旧卡——0 卡保护的窗级版)
if passed:
conn.execute(
"""UPDATE muse_knowledge_draft SET deleted=TRUE, updater=%s
WHERE tenant_id=%s AND source_type='parse_book' AND source_id=%s
AND draft_payload->>'窗起'=%s AND deleted=FALSE""",
(ACTOR, TENANT, work_id, str(from_order)))
sess = _session() # 嵌入通道(判重+新卡即时嵌入共用)
ok, merged = 0, []
for i, payload, inss in passed:
# ── 受控点①②:嵌入判重(失败不阻断——判重是增强不是红线,卡不能因通道故障丢)──
embed_text = build_embed_text(payload)
vec = None
try:
vecs, bad = embed_texts(sess, [embed_text])
vec = None if bad else vecs[0]
except Exception:
pass
if vec is not None:
top = conn.execute(
"""SELECT d.id, 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='parse_book'
ORDER BY score DESC LIMIT 1""", (json.dumps(vec), TENANT)).fetchone()
# 前向修复(2026-07-17):不再"只有 ≥SIM_MERGE 才判、0.75-0.85 仅打标记"——
# 最近邻 ≥SIM_MARK(0.75) 即送 merge_judge,与 recluster 回溯聚类同一判据。
# WHY:跨书 0.75-0.85 的近义卡此前从不真正判定、只挂"未达判定线",是重复卡沉积的病根。
if top and top[2] >= SIM_MARK:
old_id, old_payload, score = top
verdict, why = merge_judge(payload, old_payload, model)
if verdict == "merge":
# 归并=旧卡追加实例(跨书时带书名)+完整审计(含被归并卡全文,可逆)
add = [dict(ins, 书=book) for ins in inss]
old_payload["实例"] = (old_payload.get("实例") or []) + add
old_payload.setdefault("归并审计", []).append(
{"相似度": round(float(score), 4), "判定依据": why, "被归并卡": payload})
conn.execute(
"UPDATE muse_knowledge_draft SET draft_payload=%s, updater=%s WHERE id=%s",
(Jsonb(old_payload), ACTOR, old_id))
merged.append(f"《{payload['名称']}》并入已有卡#{old_id}《{old_payload.get('名称')}》({score:.2f}) {why}")
continue
# keep 留痕;档位区分强候选(≥SIM_MERGE)/扩展候选(≥SIM_MARK),供审核判读近似度
payload["判重"] = {"相似卡": old_payload.get("名称"), "相似度": round(float(score), 4),
"判定": "keep", "档": "强候选" if score >= SIM_MERGE else "扩展候选",
"依据": why}
cid = f"parse-{work_id}-w{from_order}-{i}-" + hashlib.sha256(
json.dumps(payload, ensure_ascii=False, sort_keys=True).encode()).hexdigest()[:8]
row = conn.execute(
"""INSERT INTO muse_knowledge_draft (work_id, draft_type, draft_payload, status,
source_type, source_id, command_id, creator, updater, tenant_id)
VALUES (0,'entity',%s,'pending','parse_book',%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, command_id) WHERE command_id IS NOT NULL DO NOTHING
RETURNING id""",
(Jsonb(payload), work_id, cid, ACTOR, ACTOR, TENANT)).fetchone()
# 新卡即时嵌入(下一窗/下一书判重立即可见;失败留给 embed skill 批量补)
if row and vec is not None:
h = hashlib.sha256(f"{embed_text}|{EMBED_MODEL}".encode()).hexdigest()
conn.execute(
"""INSERT INTO example_knowledge_embedding
(draft_id, content_hash, embed_text, model, dimensions, embedding,
creator, updater, tenant_id)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)
ON CONFLICT (tenant_id, content_hash, model) DO NOTHING""",
(row[0], h, embed_text, EMBED_MODEL, 1024, json.dumps(vec), ACTOR, ACTOR, TENANT))
ok += 1
# 窗内章 pattern_status 批量置 done(窗级出卡完成=这些章的范式阶段完成)
conn.execute(
"""UPDATE example_parse_task SET pattern_status='done', updater=%s
WHERE tenant_id=%s AND work_id=%s AND chapter_id IN (
SELECT id FROM muse_content_chapter
WHERE tenant_id=%s AND work_id=%s AND order_no BETWEEN %s AND %s AND deleted=FALSE)""",
(ACTOR, TENANT, work_id, TENANT, work_id, from_order, to_order))
conn.commit()
click.echo(f"cards✓ work={work_id} win#{from_order}–{to_order}: 入库{ok} 归并{len(merged)} 拒{len(rejected)}")
for m in merged:
click.echo(f" [归并] {m}")
for r in rejected:
click.echo(f" [拒] {r['卡']}: {'; '.join(r['原因'])}")
def _judge_view(payload):
"""送 merge_judge 的"手法身份"视图:只留 型/名称/摘要/字段(判"是否同一手法"的全部依据)。
WHY 裁剪:审核(三角色判词)/实例/出处/判重 是噪声——尤其母卡累积实例与「回溯归并审计」后
payload 会塞进多张别的卡全文,直接喂 merge_judge 既爆 token 又把判定带偏;裁到身份四件套后
判定输入稳定(母卡吸并再多,身份不变)、便宜、聚焦手法本身。"""
return {"型": payload.get("型"), "名称": payload.get("名称"),
"一句话摘要": payload.get("一句话摘要"), "字段": payload.get("字段") or {}}
def _mother_key(payload, cid):
"""母卡优先序(best-first 排序键):审核判定优 > 均分高 > 实例多 > id 小。
WHY:让存活下来的母卡是该手法的最佳代表——下游审核/导出以母卡为规范样本;实例多者是更厚的
证据锚点;id 兜底保证确定性可复现。审核缺失(他型可能未审)排在最后,不抢占母卡位。"""
r = payload.get("审核") or {}
rank = {"pass": 0, "revise": 1, "reject": 2}.get(r.get("判定"), 3)
return (rank, -(r.get("均分") or 0), -len(payload.get("实例") or []), cid)
def _absorb(mother, loser, sim, why):
"""把输家并入母卡(内存态,写库延后到写段统一落):输家实例(带书名归属)追加进母卡实例数组
+ 记「回溯归并审计」(相似度/判定依据/被归并卡全文——全文留痕保证归并可逆)。"""
book = (loser.get("出处") or {}).get("书名") or ""
add = [dict(ins, 书=book) for ins in (loser.get("实例") or [])] # 跨书归属:每条实例标来源书
mother["实例"] = (mother.get("实例") or []) + add
mother.setdefault("回溯归并审计", []).append(
{"相似度": round(float(sim), 4), "判定依据": why, "被归并卡": loser})
def _sim_band(s):
"""附着边相似度分档(供 dry-run 报告:不同档 LLM 剪枝率差异大,同型词汇共享令 0.75-0.80 多为假近)。"""
return "≥0.85" if s >= 0.85 else ("0.80–0.85" if s >= 0.80 else "0.75–0.80")
@cli.command()
@click.option("--type", "ptype", required=True, type=click.Choice(sorted(PATTERN_TYPES)),
help="回溯聚类的范式型(一次只跑一型)")
@click.option("--dry-run", is_flag=True, help="只算不写:报簇结构/计划归并/N→M + 抽样 merge_judge 判定")
@click.option("--limit", type=int, default=0, help="只取前 N 张卡(调试用,0=全型)")
@click.option("--sample", type=int, default=10, show_default=True,
help="dry-run 下抽样真跑 merge_judge 的候选对数(按相似度档分层取,露出各档合并率)")
@click.option("--model", default="MiniMax-M3", show_default=True, help="归并判定用模型")
def recluster(ptype, dry_run, limit, sample, model):
"""回溯聚类去重:把同型近义卡(嵌入初筛 + merge_judge 终判)归并到母卡。
病根:现行 cards() 入库只比最近 1 张、阈值偏高、按书顺序入库→跨书 0.75-0.85 的近义卡从不送判定。
本命令一次回溯:同型全量按余弦相似度成簇,逐对经 LLM 确认后把输家并入母卡(软删、可逆)。
连接三段式(血泪教训:DB 事务里夹 LLM,长空转会被 Tailscale 掐断整批崩):
读段 短连接取全型卡 + 候选对(pgvector self-join,纯 SQL 无 LLM),读完即释放;
算段 本地贪心成簇 + merge_judge 逐对终判(此阶段 recluster 不持任何 DB 连接,
merge_judge 走 chat_governed 自持额度账本短连接,互不干扰);
写段 全新短连接批量落库(母卡累积 payload / 输家卡 deleted / 输家嵌入 deleted),一次提交。
双保险:embedding 只初筛出候选,是否同一手法一律由 merge_judge 定夺,拿不准 keep(宁重复不误并)。"""
# ── 读段:短连接读完即释放 ──
with psycopg.connect(DSN) as rconn:
rows = rconn.execute(
"""SELECT d.id, d.draft_payload FROM muse_knowledge_draft d
WHERE d.tenant_id=%s AND d.source_type='parse_book' AND d.deleted=FALSE
AND d.draft_payload->>'型'=%s ORDER BY d.id""",
(TENANT, ptype)).fetchall()
if limit:
rows = rows[:limit]
cards = {cid: p for cid, p in rows}
ids = set(cards)
# 候选对:同型嵌入 self-join,余弦相似度 ≥SIM_MARK。纯 SQL(无 LLM);pgvector 一次算完 O(N²),
# combat 886 张实测 ~3s。取回后本地成簇,判定/写库阶段不再回这条连接。
pair_rows = rconn.execute(
"""WITH emb AS (
SELECT e.draft_id, e.embedding 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='parse_book' AND d.draft_payload->>'型'=%s)
SELECT a.draft_id, b.draft_id, 1-(a.embedding<=>b.embedding) AS sim
FROM emb a JOIN emb b ON a.draft_id<b.draft_id
WHERE 1-(a.embedding<=>b.embedding) >= %s""",
(TENANT, ptype, SIM_MARK)).fetchall()
# 邻接表(id→[(邻卡,相似度)]);--limit 调试时过滤掉越界 id
adj = {}
for a, b, s in pair_rows:
if a in ids and b in ids:
s = float(s)
adj.setdefault(a, []).append((b, s))
adj.setdefault(b, []).append((a, s))
n = len(ids)
if not n:
click.echo(f"[recluster·{ptype}] 无活卡")
return
# ── 算段:贪心「最近母卡」聚类(best-first)──
# 为何不用 union-find:0.75 阈下同型卡词汇高度共享,单链传递会把几乎整型并成一个巨簇
#(combat 实测 38665 对 / 817 张连通)。贪心「只附着到最近的一个母卡」天然避免链式膨胀,
# 且判定次数最省:每张非母卡至多 1 次 merge_judge(判它 vs 最近母卡)。
order = sorted(ids, key=lambda cid: _mother_key(cards[cid], cid))
leaders, leader_set = [], set()
plan = [] # 计划/已确认归并:(母卡id, 输家id, 相似度)
def nearest_leader(cid):
"""该卡邻居中"已是母卡"的最近一个(sim 最大);无则 None。"""
best = None
for nbr, s in adj.get(cid, []):
if nbr in leader_set and (best is None or s > best[1]):
best = (nbr, s)
return best
if dry_run:
# 纯向量投影(几乎不烧 LLM):每卡附着到最近母卡=一条计划归并,无最近母卡=自成母卡。
# 这是归并的「上限」——实跑每对还要过 merge_judge,keep 者不并→真实存活更多。
for cid in order:
nb = nearest_leader(cid)
if nb:
plan.append((nb[0], cid, nb[1]))
else:
leaders.append(cid)
leader_set.add(cid)
else:
# 实跑:每次附着都由 merge_judge 终判(受控点②)。此阶段 recluster 不持 DB 连接。
for cid in order:
nb = nearest_leader(cid)
if nb:
lid, sim = nb
verdict, why = merge_judge(_judge_view(cards[cid]), _judge_view(cards[lid]), model)
if verdict == "merge":
_absorb(cards[lid], cards[cid], sim, why) # 内存改母卡,写段统一落库
plan.append((lid, cid, sim))
continue
leaders.append(cid) # 无最近母卡 或 判 keep → 自成母卡
leader_set.add(cid)
# ── 写段:全新短连接批量落库(红线:软删,绝不物理删)──
if not dry_run and plan:
movers = {lid for lid, _lo, _s in plan} # 被吸并过的母卡(去重,每张只写一次最终态)
with psycopg.connect(DSN) as wconn:
for lid in movers:
wconn.execute("UPDATE muse_knowledge_draft SET draft_payload=%s, updater=%s WHERE id=%s",
(Jsonb(cards[lid]), ACTOR, lid))
for _lid, loser_id, _s in plan:
wconn.execute("UPDATE muse_knowledge_draft SET deleted=TRUE, updater=%s WHERE id=%s",
(ACTOR, loser_id)) # 输家卡软删
wconn.execute(
"""UPDATE example_knowledge_embedding SET deleted=TRUE, updater=%s
WHERE tenant_id=%s AND draft_id=%s AND deleted=FALSE""",
(ACTOR, TENANT, loser_id)) # 输家嵌入软删:未来判重/检索不再命中已并走的卡
wconn.commit()
# ── 报告 ──
from collections import Counter
per_mother = Counter(lid for lid, _lo, _s in plan) # 母卡 → 吸并数
band = Counter(_sim_band(s) for _l, _lo, s in plan) # 附着边相似度分档
surv = n - len(plan) # 存活母卡数
tag = "dry-run·向量投影上限" if dry_run else "实跑·LLM确认"
click.echo(f"[recluster·{ptype}·{tag}] 活卡 {n} → 存活 {surv}"
f"(归并 {len(plan)},缩减 {len(plan) / n:.0%})")
if dry_run:
click.echo(" 注:向量投影是归并上限;实跑每对经 merge_judge 终判,keep 者不并→实际存活更多。")
click.echo(" 附着边相似度:" + " ".join(f"{k}={band.get(k, 0)}" for k in ("≥0.85", "0.80–0.85", "0.75–0.80")))
click.echo(f" 含归并的簇 {len(per_mother)}(另有 {surv - len(per_mother)} 张孤卡母卡);最大簇 top5:")
for mid, cnt in per_mother.most_common(5):
click.echo(f" #{mid}《{cards[mid].get('名称')}》聚 {cnt + 1} 张(母1+并{cnt})")
# dry-run 抽样:按相似度档「分层」各取若干真跑 merge_judge——展示判定质量 + 分档外推真实归并数。
# WHY 分层而非等距:附着边多数落在 0.75-0.80(同型词汇高度共享的"假近"区,LLM 几乎全 keep),
# 等距抽样会几乎全落该档、把高档才有的真归并淹没成 0;分层能露出每档合并率,据此把向量上限收敛到现实预估。
if dry_run and plan and sample > 0:
by_band = {"≥0.85": [], "0.80–0.85": [], "0.75–0.80": []}
for e in sorted(plan, key=lambda x: x[2], reverse=True):
by_band[_sim_band(e[2])].append(e)
live_bands = [b for b in by_band if by_band[b]]
per = max(1, sample // max(1, len(live_bands))) # 每档配额(档内等距取,覆盖该档相似度跨度)
picks = []
for b in live_bands:
es = by_band[b]
k = min(per, len(es))
idxs = [round(i * (len(es) - 1) / (k - 1)) for i in range(k)] if k > 1 else [0]
picks += [es[i] for i in dict.fromkeys(idxs)]
click.echo(f" —— 抽样 merge_judge 判定({len(picks)} 对,按相似度档分层)——")
band_hit = {b: [0, 0] for b in by_band} # 档 → [merge数, 抽样数]
for lid, loser_id, sim in picks:
verdict, why = merge_judge(_judge_view(cards[loser_id]), _judge_view(cards[lid]), model)
bd = _sim_band(sim)
band_hit[bd][0] += verdict == "merge"
band_hit[bd][1] += 1
click.echo(f" [{verdict}] {sim:.3f}({bd}) 输《{cards[loser_id].get('名称')}》→ "
f"母《{cards[lid].get('名称')}》:{why}")
# 分档外推:估真实归并 = Σ(该档附着边数 × 该档抽样 merge 率)。样本小,仅供数量级判断。
est = sum(band.get(b, 0) * (band_hit[b][0] / band_hit[b][1] if band_hit[b][1] else 0) for b in by_band)
for b in by_band:
hit, tot = band_hit[b]
click.echo(f" 档 {b}: 抽样 merge {hit}/{tot},本档附着边 {band.get(b, 0)}"
f" → 估归并 ~{band.get(b, 0) * (hit / tot if tot else 0):.0f}")
click.echo(f" 分档外推真实归并 ~{est:.0f}({n}→~{n - est:.0f});样本小仅供数量级,放量以实跑逐对确认为准")
elif not dry_run:
for lid, loser_id, _s in plan[:20]:
click.echo(f" [归并] 《{cards[loser_id].get('名称')}》并入 #{lid}《{cards[lid].get('名称')}》")
if len(plan) > 20:
click.echo(f" …另 {len(plan) - 20} 条归并")
@cli.command()
@click.option("--work-id", type=int)
def progress(work_id):
"""进度统计(审查面)。"""
with psycopg.connect(DSN) as conn:
where = " AND t.work_id=%s" if work_id else ""
args = [TENANT] + ([work_id] if work_id else [])
rows = conn.execute(f"""
SELECT w.title, count(*) FILTER (WHERE t.scaffold_status='done') AS s_done,
count(*) FILTER (WHERE t.pattern_status='done') AS p_done, count(*) AS total,
(SELECT count(*) FROM muse_knowledge_draft d
WHERE d.tenant_id=%s AND d.source_type='parse_book' AND d.source_id=t.work_id
AND d.deleted=FALSE) AS drafts
FROM example_parse_task t JOIN muse_content_work w ON w.id=t.work_id
WHERE t.tenant_id=%s AND t.deleted=FALSE{where}
GROUP BY w.title, t.work_id ORDER BY w.title""", [TENANT] + args).fetchall()
for r in rows:
click.echo(f"{r[0]:<12} 脚手架 {r[1]}/{r[3]} 范式 {r[2]}/{r[3]} 草稿卡 {r[4]}")
if __name__ == "__main__":
try:
cli()
except psycopg.Error as e:
click.echo(f"[db错误] {type(e).__name__}: {e}", err=True)
sys.exit(1)