533 lines
32 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, 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.85 # 嵌入判重:≥该值触发 M3 归并判定(阈值未校准,实战收集中)
SIM_MARK = 0.75 # [MARK, MERGE) 区间只标记不判定(校准带宽,审核环可见)
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-S2craft 双模板条件必填/禁填 + 装置类型闭合枚举 + 禁复合标签
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: # 拍板值 35%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, _ = chat(prompt, model=model)
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()
if top and top[2] >= SIM_MERGE:
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
payload["判重"] = {"相似卡": old_payload.get("名称"), "相似度": round(float(score), 4),
"判定": "keep", "依据": why}
elif top and top[2] >= SIM_MARK:
payload["判重"] = {"相似卡": top[1].get("名称"), "相似度": round(float(top[2]), 4),
"判定": "未达判定线"}
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['原因'])}")
@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)