#!/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 背景词不在实体名池)。 # 注意"浮游炮"不入表——已泛化为中文机甲文通用武器品类词(同"光剑"),入表实测误伤 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-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, _ = 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)