249 lines
12 KiB
Python
249 lines
12 KiB
Python
#!/usr/bin/env python3
|
||
"""parse-book 配套确定性脚本:拆书产物校验入库 + 任务状态机(B2)。
|
||
|
||
职责边界:extractor(LLM) 只产结构化 JSON 文件,不碰库;本脚本做机械校验后写库——
|
||
字段 key 合法性(对库内字段合同)、五型归型、出处必填、**脱敏红线 15 连字检测**(硬阻断)。
|
||
状态全在库(example_parse_task),断点续跑与幂等按章。
|
||
"""
|
||
import hashlib
|
||
import json
|
||
import pathlib
|
||
import re
|
||
import sys
|
||
|
||
import click
|
||
import psycopg
|
||
from psycopg.types.json import Jsonb
|
||
|
||
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)
|
||
|
||
|
||
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")})
|
||
return {"细纲": data.get("细纲") or data.get("outline") or "", "实体": ents}
|
||
|
||
|
||
def norm_card(c: dict) -> dict:
|
||
"""范式卡键名归一(ASCII→中文)。"""
|
||
src = c.get("出处") or c.get("source") 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 {},
|
||
"出处": {"书名": 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 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 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_):
|
||
"""脚手架入库:{细纲, 实体:[{型,名称,一句话摘要,备注?}]};比例约束校验(3–5%,超标拒绝)。"""
|
||
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)))
|
||
if ratio > 0.08: # 拍板值 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%),已记失败退回重解析")
|
||
ents = data.get("实体") or []
|
||
for e in ents:
|
||
if not e.get("名称") or not e.get("型"):
|
||
raise click.ClickException(f"实体缺 名称/型: {e}")
|
||
conn.execute(
|
||
"""INSERT INTO example_parse_scaffold (work_id, chapter_id, outline_text, entities,
|
||
creator, updater, tenant_id)
|
||
VALUES (%s,%s,%s,%s,%s,%s,%s)
|
||
ON CONFLICT (tenant_id, chapter_id)
|
||
DO UPDATE SET outline_text=EXCLUDED.outline_text, entities=EXCLUDED.entities,
|
||
deleted=FALSE, updater=EXCLUDED.updater""",
|
||
(work_id, ch_id, outline, Jsonb(ents), 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%}) 实体{len(ents)}")
|
||
|
||
|
||
@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_):
|
||
"""范式卡入库:[{型∈五型, 名称, 一句话摘要, 字段{…}, 出处{书名,回目,定位}}] → draft(pending)。
|
||
机械硬阻断:归型合法、字段 key 合法(库内合同)、出处必填、15 连字脱敏检测。"""
|
||
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]
|
||
# 幂等:重跑本章 = 软删本章旧 draft
|
||
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):
|
||
t = c.get("型")
|
||
reasons = []
|
||
if t not in PATTERN_TYPES:
|
||
reasons.append(f"型不合法:{t}(首轮只拆五型)")
|
||
if not c.get("名称"):
|
||
reasons.append("缺名称")
|
||
src_ref = c.get("出处") or {}
|
||
if not (src_ref.get("书名") and src_ref.get("回目")):
|
||
reasons.append("出处不完整(需书名+回目)")
|
||
fields = c.get("字段") or {}
|
||
if t in contracts:
|
||
illegal = set(fields) - contracts[t]
|
||
if illegal:
|
||
reasons.append(f"字段越合同:{sorted(illegal)}")
|
||
texts = [c.get("名称"), c.get("一句话摘要"), *fields.values()]
|
||
leak = leak_check([x for x in texts if x], src)
|
||
if leak:
|
||
reasons.append(f"脱敏违规(≥{NGRAM}连字重合):「{leak}」")
|
||
if reasons:
|
||
rejected.append({"卡": c.get("名称") or f"#{i}", "原因": reasons})
|
||
continue
|
||
payload = {"型": t, "名称": c["名称"], "一句话摘要": c.get("一句话摘要", ""),
|
||
"字段": fields, "出处": src_ref, "目标库": "公共范式库",
|
||
"章序": chapter_order, "来源": f"拆书@{book}", "状态": "草稿"}
|
||
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" if not rejected else "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['原因'])}")
|
||
|
||
|
||
@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)
|