#!/usr/bin/env python3 """从已接受章节抽取作品知识草稿。 本入口只负责“正文 -> draft/shadow”这一段,不自动确认知识;确认由 decide-candidate 单独触发。模型调用经 call-content-model,所有输入输出挂到本次 run_id。 """ import hashlib import json import pathlib import sys from typing import Any import click import psycopg from psycopg.types.json import Jsonb HERE = pathlib.Path(__file__).resolve().parent SKILLS = HERE.parents[1] for import_path in ( SKILLS / "call-content-model" / "scripts", SKILLS / "record-run-evidence" / "scripts", ): if str(import_path) not in sys.path: sys.path.insert(0, str(import_path)) from llm import chat_governed, cost_usd, extract_json # noqa: E402 from record_failed_run import record_failure # noqa: E402 from run_registry import finish_run, new_run_id, start_run # noqa: E402 DB_SCRIPTS = SKILLS / "access-database" / "scripts" if str(DB_SCRIPTS) not in sys.path: sys.path.insert(0, str(DB_SCRIPTS)) from db import connect # noqa: E402 TENANT, ACTOR = 1, "1" MODEL = "MiniMax-M3" ENTITY_TYPES = frozenset({ "character", "location", "faction", "power_system", "item", "event", }) TYPE_ALIASES = { "人物": "character", "地点": "location", "组织": "faction", "势力": "faction", "能力体系": "power_system", "力量体系": "power_system", "物件": "item", "事件": "event", } class ExtractionContractError(ValueError): """模型输出无法绑定到正文事实时失败关闭。""" PROMPT = """你是长篇小说章后知识抽取员。只从给定的已接受正文抽取作品私有知识草稿。 不要确认知识,不要补写正文没有的事实;低置信内容仍保留但在 brief/fields 中标注“?”。 实体类型只能使用:character、location、faction、power_system、item、event。 立卡门槛:具名且有跨章复用或后续履约潜力;一次性龙套和一次性道具不要列实体。 证据必须是正文中的逐字连续片段,不能改写。 只输出一个 JSON 对象,严格使用以下 ASCII 字段,不要 markdown: {{ "entities": [{{"type":"character", "name":"", "brief":"", "fields": {{}}, "evidence":"正文逐字片段"}}], "relations": [{{"source":"实体名", "target":"实体名", "type":"关系类型", "description":"", "evidence":"正文逐字片段"}}], "state": {{"currentSituation":"", "characterStates": {{}}, "foreshadowing": {{"埋":[],"推":[],"收":[]}}, "handoff":""}} }} 作品:《{title}》 章节:第 {chapter_order} 章《{chapter_title}》 已有确认实体名(只用于判重,不得把没有正文证据的内容写进本章):{existing_names} 【正文】 {body} """ def _first(item: dict[str, Any], *keys, default=None): for key in keys: if key in item: return item[key] return default def _text(value, field): if not isinstance(value, str) or not value.strip(): raise ExtractionContractError(f"{field} 必须是非空字符串") return value.strip() def _hash_payload(value: Any) -> str: return hashlib.sha256( json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8") ).hexdigest() def _merge_usage(total, current): """只累加顶层数值 token,保留嵌套明细的简单形状。""" for key, value in (current or {}).items(): if isinstance(value, (int, float)): total[key] = total.get(key, 0) + value def normalize_extraction(raw: Any, body: str) -> dict[str, Any]: """把模型输出归一化并机械绑定到本章正文。""" if not isinstance(raw, dict): raise ExtractionContractError("抽取输出必须是对象") entities = raw.get("entities") or raw.get("实体") or [] relations = raw.get("relations") or raw.get("关系") or [] state = raw.get("state") or raw.get("状态") or {} if not isinstance(entities, list) or not isinstance(relations, list) or not isinstance(state, dict): raise ExtractionContractError("entities/relations/state 类型非法") normalized_entities = [] names = set() for index, item in enumerate(entities): if not isinstance(item, dict): raise ExtractionContractError(f"entities[{index}] 必须是对象") entity_type = _text(_first(item, "type", "型"), f"entities[{index}].type") entity_type = TYPE_ALIASES.get(entity_type, entity_type) if entity_type not in ENTITY_TYPES: raise ExtractionContractError(f"entities[{index}].type 非法: {entity_type}") name = _text(_first(item, "name", "名称"), f"entities[{index}].name") brief = _text(_first(item, "brief", "一句话摘要", default="?"), f"entities[{index}].brief") evidence = _text(_first(item, "evidence", "证据"), f"entities[{index}].evidence") if evidence not in body: raise ExtractionContractError(f"entities[{index}] 证据不在正文中") key = (entity_type, name.casefold()) if key in names: raise ExtractionContractError(f"实体重复: {entity_type}/{name}") names.add(key) fields = _first(item, "fields", "字段", default={}) if not isinstance(fields, dict): raise ExtractionContractError(f"entities[{index}].fields 必须是对象") normalized_entities.append({ "type": entity_type, "name": name, "brief": brief, "fields": fields, "evidence": evidence, }) normalized_relations = [] for index, item in enumerate(relations): if not isinstance(item, dict): raise ExtractionContractError(f"relations[{index}] 必须是对象") source = _text(_first(item, "source", "甲方"), f"relations[{index}].source") target = _text(_first(item, "target", "乙方"), f"relations[{index}].target") relation_type = _text(_first(item, "type", "关系类型"), f"relations[{index}].type") description = _text(_first(item, "description", "描述", default="?"), f"relations[{index}].description") evidence = _text(_first(item, "evidence", "证据"), f"relations[{index}].evidence") if evidence not in body: raise ExtractionContractError(f"relations[{index}] 证据不在正文中") normalized_relations.append({ "source": source, "target": target, "type": relation_type, "description": description, "evidence": evidence, }) return { "entities": normalized_entities, "relations": normalized_relations, "state": state, } def salvage_extraction(raw: Any, body: str) -> dict[str, Any]: """删除无法绑定的模型条目,再复用同一严格归一器。 这是保守收口,不替模型编造证据:实体名本身若逐字出现在正文,可作为最小证据; 关系缺证据则直接丢弃。丢弃数量写入 payload,供质量结果和看板解释。 """ if not isinstance(raw, dict): raise ExtractionContractError("无法从非对象输出做保守收口") candidate = dict(raw) entity_key = "entities" if "entities" in candidate else "实体" relation_key = "relations" if "relations" in candidate else "关系" kept_entities, dropped_entities = [], 0 for item in candidate.get(entity_key) or []: if not isinstance(item, dict): dropped_entities += 1 continue name = _first(item, "name", "名称") evidence = _first(item, "evidence", "证据") if isinstance(evidence, str) and evidence in body: kept_entities.append(item) elif isinstance(name, str) and name.strip() and name.strip() in body: fixed = dict(item) fixed["evidence"] = name.strip() kept_entities.append(fixed) else: dropped_entities += 1 kept_relations, dropped_relations = [], 0 for item in candidate.get(relation_key) or []: if not isinstance(item, dict) or not isinstance(_first(item, "evidence", "证据"), str) \ or _first(item, "evidence", "证据") not in body: dropped_relations += 1 continue kept_relations.append(item) candidate[entity_key] = kept_entities candidate[relation_key] = kept_relations normalized = normalize_extraction(candidate, body) normalized["mechanicalDrops"] = { "entities": dropped_entities, "relations": dropped_relations, } return normalized def _load_chapter(work_id, chapter_order): with connect(readonly=True) as conn: row = conn.execute( "SELECT w.title,c.id,c.title,b.id,b.content_text " "FROM muse_content_work w JOIN muse_content_chapter c ON c.work_id=w.id " "JOIN muse_content_block b ON b.chapter_id=c.id AND b.deleted=false " "WHERE w.id=%s AND c.order_no=%s AND w.deleted=false AND c.deleted=false " "ORDER BY b.revision DESC LIMIT 1", (work_id, chapter_order), ).fetchone() if not row: raise ValueError(f"作品 {work_id} 第 {chapter_order} 章没有可用 Canonical 正文") existing = [r[0] for r in conn.execute( "SELECT normalized_name FROM muse_knowledge_entity WHERE tenant_id=%s AND work_id=%s " "AND deleted=false ORDER BY id", (TENANT, work_id) ).fetchall()] return row, existing def _insert_draft(conn, *, work_id, chapter_id, chapter_order, run_id, entity, index): payload = { "type": entity["type"], "name": entity["name"], "brief": entity["brief"], "fields": entity["fields"], "evidence": entity["evidence"], "source": {"workId": work_id, "chapter": chapter_order, "chapterId": chapter_id}, "extractRunId": run_id, } command_id = f"extract-{work_id}-ch{chapter_order}-entity-{index}-{_hash_payload(payload)[:12]}" normalized_name = entity["name"].strip().casefold() current = conn.execute( "SELECT id,description,attributes,revision FROM muse_knowledge_entity " "WHERE tenant_id=%s AND work_id=%s AND entity_type=%s AND normalized_name=%s " "AND scope='local' AND deleted=false FOR SHARE", (TENANT, work_id, entity["type"], normalized_name), ).fetchone() existing_id = current[0] if current else None snapshot = None if not current else { "id": current[0], "description": current[1], "attributes": current[2], "revision": current[3] } row = conn.execute( "INSERT INTO muse_knowledge_draft(work_id,entity_id,draft_type,target_object_id,proposed_changes," "current_canonical_snapshot,draft_payload,source_status,source_action_policy,status,confidence," "source_type,source_id,command_id,creator,updater,tenant_id) " "VALUES (%s,%s,'entity',%s,%s::jsonb,%s::jsonb,%s::jsonb,'active','allowed','pending',%s," "'chapter_extract',%s,%s,%s,%s,%s) " "ON CONFLICT (tenant_id,command_id) WHERE command_id IS NOT NULL DO NOTHING RETURNING id", ( work_id, existing_id, existing_id, json.dumps({"description": entity["brief"], "attributes": entity["fields"]}, ensure_ascii=False), json.dumps(snapshot, ensure_ascii=False) if snapshot else None, json.dumps(payload, ensure_ascii=False), 0.8, chapter_id, command_id, ACTOR, ACTOR, TENANT, ), ).fetchone() return row[0] if row else None def persist_extraction(work_id, chapter_id, chapter_order, run_id, payload, *, requested_model, actual_model, usage): """一次事务写实体/关系草稿、状态 shadow、运行回执和质量结果。""" result_sha = _hash_payload(payload) with connect() as conn: try: draft_ids = [] for index, entity in enumerate(payload["entities"], start=1): draft_id = _insert_draft( conn, work_id=work_id, chapter_id=chapter_id, chapter_order=chapter_order, run_id=run_id, entity=entity, index=index, ) if draft_id: draft_ids.append(draft_id) for index, relation in enumerate(payload["relations"], start=1): relation_payload = { **relation, "source": {"name": relation["source"]}, "target": {"name": relation["target"]}, "sourceRef": {"workId": work_id, "chapter": chapter_order, "chapterId": chapter_id}, "extractRunId": run_id, } command_id = f"extract-{work_id}-ch{chapter_order}-relation-{index}-{_hash_payload(relation_payload)[:12]}" row = conn.execute( "INSERT INTO muse_knowledge_draft(work_id,draft_type,draft_payload,source_status," "source_action_policy,status,confidence,source_type,source_id,command_id,creator,updater,tenant_id) " "VALUES (%s,'relation',%s::jsonb,'active','allowed','pending',%s,'chapter_extract',%s,%s,%s,%s,%s) " "ON CONFLICT (tenant_id,command_id) WHERE command_id IS NOT NULL DO NOTHING RETURNING id", (work_id, json.dumps(relation_payload, ensure_ascii=False), 0.7, chapter_id, command_id, ACTOR, ACTOR, TENANT), ).fetchone() if row: draft_ids.append(row[0]) state = payload.get("state") or {} state_payload = { "schemaVersion": "narrative-state-v1", "workId": work_id, "chapter": chapter_order, "state": state, "source": {"chapterId": chapter_id, "runId": run_id}, } state_id = None if state: version = conn.execute( "SELECT COALESCE(MAX(version),0)+1 FROM example_planning_section " "WHERE tenant_id=%s AND work_id=%s AND section_type='state' AND target_chapter IS NULL", (TENANT, work_id), ).fetchone()[0] state_id = conn.execute( "INSERT INTO example_planning_section(work_id,section_type,schema_type,version,payload,state,creator,updater,tenant_id) " "VALUES (%s,'state','narrative_state',%s,%s::jsonb,'shadow',%s,%s,%s) RETURNING id", (work_id, version, json.dumps(state_payload, ensure_ascii=False), ACTOR, ACTOR, TENANT), ).fetchone()[0] raw_id = conn.execute( "SELECT raw_content_id FROM example_llm_call WHERE run_id=%s AND caller='extract-knowledge' " "ORDER BY id DESC LIMIT 1", (run_id,) ).fetchone() raw_content_id = raw_id[0] if raw_id else None receipt_id = conn.execute( "INSERT INTO example_run_receipt(run_id,sample_id,revision,adapter_role,stage_kind,attempt," "requested_model_id,actual_model_id,model_match,total_cost_usd,usage,stop_reason,terminal_reason," "is_error,safe_summary,result_sha256,raw_content_id,creator,tenant_id) " "VALUES (%s,%s,1,'extractor','generation',1,%s,%s,%s,%s,%s::jsonb,'stop','completed',FALSE,%s::jsonb,%s,%s,%s,%s) " "ON CONFLICT (tenant_id,run_id,sample_id,revision) DO NOTHING RETURNING id", ( run_id, f"extract-ch{chapter_order}", requested_model, actual_model, requested_model == actual_model, cost_usd(actual_model, usage), json.dumps(usage, ensure_ascii=False), json.dumps({"entityDrafts": len(draft_ids), "stateDraftId": state_id, "mechanicalDrops": payload.get("mechanicalDrops", {})}, ensure_ascii=False), result_sha, raw_content_id, ACTOR, TENANT, ), ).fetchone() receipt_id = receipt_id[0] if receipt_id else None if receipt_id is None: receipt_id = conn.execute( "SELECT id FROM example_run_receipt WHERE tenant_id=%s AND run_id=%s AND sample_id=%s AND revision=1", (TENANT, run_id, f"extract-ch{chapter_order}"), ).fetchone()[0] conn.execute( "INSERT INTO example_quality_result(run_id,receipt_id,judge_kind,scale_version,conclusion,detail,raw_content_id,creator,tenant_id) " "VALUES (%s,%s,'detection','extractor-contract-v1','pass',%s::jsonb,%s,%s,%s) " "ON CONFLICT (tenant_id,run_id,judge_kind,COALESCE(dimension,''),COALESCE(candidate_sha256,'')) DO NOTHING", (run_id, receipt_id, json.dumps({"entityDrafts": len(draft_ids), "relationDrafts": len(payload["relations"]), "stateDraftId": state_id, "mechanicalDrops": payload.get("mechanicalDrops", {})}, ensure_ascii=False), raw_content_id, ACTOR, TENANT), ) conn.commit() return {"draft_ids": draft_ids, "state_draft_id": state_id, "receipt_id": receipt_id, "result_sha256": result_sha} except Exception: conn.rollback() raise def extract_chapter(work_id, chapter_order, *, run_id=None): record = start_run( run_id=run_id or new_run_id("extract-knowledge", work_id=work_id, target_chapter=chapter_order), work_id=work_id, target_chapter=chapter_order, trigger_detail={"stage": "chapter-after-extraction"}, creator=ACTOR, ) active_run = record["run_id"] try: (title, chapter_id, chapter_title, _block_id, body), existing = _load_chapter(work_id, chapter_order) prompt = PROMPT.format( title=title, chapter_order=chapter_order, chapter_title=chapter_title or "", existing_names="、".join(existing[:200]) or "(暂无)", body=body, ) content, usage, actual_model = chat_governed( prompt, model=MODEL, caller="extract-knowledge", run_id=active_run, ) if actual_model is None: raise RuntimeError("抽取模型治理链全部耗尽") total_usage = dict(usage or {}) raw_output = extract_json(content) try: payload = normalize_extraction(raw_output, body) except ExtractionContractError as first_error: # 只允许模型按正文逐字重绑证据;不能借 repair 轮新增实体、关系或事实。 repair_prompt = ( prompt + "\n\n【机械校验失败,允许一次修复】\n" + f"失败原因:{first_error}\n" + "只修正证据字段,使每条 evidence 都是上方正文中的逐字连续片段;" "删除无法找到逐字证据的条目,不得新增条目、事实、关系或状态。仍只输出同一 JSON 对象。" ) content, repair_usage, repair_model = chat_governed( repair_prompt, model=MODEL, caller="extract-knowledge", run_id=active_run, ) if repair_model is None: raise _merge_usage(total_usage, repair_usage) repaired_raw = extract_json(content) try: payload = normalize_extraction(repaired_raw, body) except ExtractionContractError: payload = salvage_extraction(repaired_raw, body) actual_model = repair_model result = persist_extraction( work_id, chapter_id, chapter_order, active_run, payload, requested_model=MODEL, actual_model=actual_model, usage=total_usage, ) finish_run(active_run, "completed", creator=ACTOR, trigger_detail={"stage": "chapter-after-extraction", "drafts": len(result["draft_ids"])}) return {"run_id": active_run, **result, "actual_model": actual_model} except BaseException as exc: finish_run(active_run, "failed", creator=ACTOR, trigger_detail={"stage": "chapter-after-extraction", "error_type": type(exc).__name__}) try: record_failure( active_run, sample_id=f"extract-ch{chapter_order}", adapter_role="extractor", caller="extract-knowledge", failure_type=type(exc).__name__, ) except Exception as receipt_error: click.echo( f"[警告] 失败回执补写失败:{type(receipt_error).__name__}", err=True, ) raise @click.command() @click.option("--work-id", type=int, required=True) @click.option("--chapter-order", type=int, required=True) @click.option("--run-id", default=None) def main(work_id, chapter_order, run_id): try: click.echo(json.dumps(extract_chapter(work_id, chapter_order, run_id=run_id), ensure_ascii=False)) except (ExtractionContractError, RuntimeError, psycopg.Error) as exc: raise click.ClickException(str(exc)) from exc if __name__ == "__main__": main()