zizi 2666d50a8e 修复: 收紧授权用途与导入任务匹配
阻断允许用途与禁止用途冲突,并让授权快照外层与内层保持一致。
导入任务按参考作品原文件唯一匹配,避免无关成功任务误阻断回放。
2026-07-19 22:19:47 +08:00

750 lines
30 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
"""从实验库只读组装回放评测配置。
本适配器只做 SELECT 和临时文件输出,不写数据库、不读取正文 block 的内容。
参考作品的目标章只进入审计侧 proxy,历史规划上下文和三臂卡注入严格分开。
"""
from __future__ import annotations
import argparse
import copy
import json
import re
from datetime import date, datetime
from pathlib import Path
from typing import Any, Mapping, Sequence
import psycopg
from psycopg.rows import dict_row
from build_snapshot import filter_milestones, filter_outline_windows, normalize_chapter
DSN = (
"postgresql://root:f6710e2d0294eb1c10e26a805a64bc54@100.64.0.8:5433/muse-example"
"?keepalives=1&keepalives_idle=15&keepalives_interval=5&keepalives_count=3"
)
TENANT_ID = 1
REPO_ROOT = Path(__file__).resolve().parents[4]
DEFAULT_SNAPSHOT_VERSION = "next_fine_outline_replay_v0"
FILE_HASH_PATTERN = re.compile(r"^[0-9a-f]{64}$")
class AdapterError(ValueError):
"""只读适配输入缺失、越界或不能证明安全时抛出。"""
def _safe_json(value: Any) -> str:
"""用固定格式序列化配置,保证运行版本可复现。"""
return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
def _required_chapter(value: Any, field: str) -> int:
"""把章号限制为明确正整数,拒绝猜测性转换。"""
chapter = normalize_chapter(value)
if chapter is None:
raise AdapterError(f"{field} 必须是正整数章号")
return chapter
def _unique_record(rows: Sequence[Mapping[str, Any]], label: str) -> Mapping[str, Any]:
"""来源证明链要求唯一行;缺失或重复都不能猜测选取。"""
if not isinstance(rows, Sequence) or isinstance(rows, (str, bytes)) or len(rows) != 1:
count = len(rows) if isinstance(rows, Sequence) and not isinstance(rows, (str, bytes)) else 0
raise AdapterError(f"{label} 必须唯一匹配,实际 {count} 行")
row = rows[0]
if not isinstance(row, Mapping):
raise AdapterError(f"{label} 不是对象")
if row.get("deleted") is True:
raise AdapterError(f"{label} 已软删")
return row
def validate_source_records(
reference_rows: Sequence[Mapping[str, Any]],
import_task_rows: Sequence[Mapping[str, Any]],
document_rows: Sequence[Mapping[str, Any]],
) -> dict[str, Any]:
"""交叉核验原文件登记、成功导入任务和知识文档,生成稳定文件版本。"""
reference = _unique_record(reference_rows, "reference.source_file")
import_task = _unique_record(import_task_rows, "成功 import task")
document = _unique_record(document_rows, "未删 knowledge document")
source_file = str(reference.get("source_file") or "").strip()
if not source_file:
raise AdapterError("reference.source_file 为空")
if str(import_task.get("status") or "") != "succeeded":
raise AdapterError("import task 不是 succeeded")
source_snapshot = import_task.get("source_snapshot")
if not isinstance(source_snapshot, Mapping):
raise AdapterError("import task 缺少 source_snapshot")
if str(source_snapshot.get("file") or "") != source_file:
raise AdapterError("import task filename 与 reference.source_file 不一致")
if str(document.get("file_name") or "") != source_file:
raise AdapterError("knowledge document filename 与 reference.source_file 不一致")
file_hash = str(document.get("file_hash") or "")
if not FILE_HASH_PATTERN.fullmatch(file_hash):
raise AdapterError("knowledge document file_hash 必须是 64 位小写十六进制")
expected_command_prefix = f"import-{file_hash[:16]}"
command_id = str(import_task.get("command_id") or "")
if not command_id.startswith(expected_command_prefix):
raise AdapterError("import task command_id 前缀与原文件 hash 不一致")
source_hash = f"sha256:{file_hash}"
return {
"referenceWorkId": str(reference.get("id") or ""),
"importTaskId": str(import_task.get("id") or ""),
"documentId": str(document.get("id") or ""),
"fileName": source_file,
"sourceHash": source_hash,
"sourceVersion": f"raw-file-v1:{source_hash}",
"importCommandId": command_id,
}
def _iso_time(value: Any) -> str | None:
"""把数据库时间统一为带时区的 ISO 字符串,空值保持为空。"""
if value is None:
return None
if isinstance(value, (datetime, date)):
return value.isoformat()
text = str(value).strip()
return text or None
def project_authorization(
row: Mapping[str, Any] | None,
source: Mapping[str, Any],
) -> dict[str, Any]:
"""把最新授权行投影为外层授权与不可变快照;无记录时保持阻断。"""
source_hash = str(source.get("sourceHash") or "")
source_version = str(source.get("sourceVersion") or "")
if row is None:
return {
"sourceStatus": "missing_authorization_snapshot",
"copyrightStatus": "unknown",
"sourceHash": source_hash,
"sourceVersion": source_version,
"allowedPurpose": [],
"forbiddenPurpose": [],
"authorizationSnapshot": {},
}
if not isinstance(row, Mapping):
raise AdapterError("授权快照行不是对象")
if str(row.get("source_hash") or "") != source_hash:
raise AdapterError("授权快照 source_hash 与原文件不一致")
if str(row.get("source_version") or "") != source_version:
raise AdapterError("授权快照 source_version 与原文件不一致")
allowed = row.get("allowed_purpose")
forbidden = row.get("forbidden_purpose")
if not isinstance(allowed, list) or not isinstance(forbidden, list):
raise AdapterError("授权快照用途字段必须是数组")
source_status = str(row.get("source_status") or "")
copyright_status = str(row.get("copyright_status") or "")
snapshot = {
"id": str(row.get("id") or ""),
"version": str(row.get("snapshot_version") or ""),
"immutable": True,
"sourceHash": source_hash,
"sourceVersion": source_version,
"sourceStatus": source_status,
"copyrightStatus": copyright_status,
"allowedPurpose": copy.deepcopy(allowed),
"forbiddenPurpose": copy.deepcopy(forbidden),
"authorizationBasis": str(row.get("authorization_basis") or ""),
"authorizedBy": str(row.get("authorized_by") or ""),
"displaySummary": str(row.get("display_summary") or ""),
"checkedAt": _iso_time(row.get("checked_at")),
"expiresAt": _iso_time(row.get("expires_at")),
"revalidationAt": _iso_time(row.get("revalidation_at")),
}
return {
"sourceStatus": source_status,
"copyrightStatus": copyright_status,
"sourceHash": source_hash,
"sourceVersion": source_version,
"allowedPurpose": copy.deepcopy(allowed),
"forbiddenPurpose": copy.deepcopy(forbidden),
"authorizationSnapshot": snapshot,
}
def _row_id(row: Mapping[str, Any]) -> str:
"""读取卡行主键并统一为来源 ID 字符串。"""
value = row.get("id")
if value is None or str(value).strip() == "":
raise AdapterError("卡行缺少 id")
return str(value)
def _normalize_id_list(value: Any, field: str) -> list[str]:
"""校验预注册卡 ID 列表,不根据目标章临时猜卡。"""
if not isinstance(value, list) or not value:
raise AdapterError(f"{field} 必须是非空数组")
if any(
isinstance(item, bool)
or not isinstance(item, (int, str))
or not str(item).strip().isdigit()
or int(item) <= 0
for item in value
):
raise AdapterError(f"{field} 只能包含正整数卡 ID")
result = [str(item) for item in value if str(item).strip()]
if len(result) != len(value) or len(result) != len(set(result)):
raise AdapterError(f"{field} 含空 ID 或重复 ID")
return result
def _card_history(payload: Mapping[str, Any]) -> list[Mapping[str, Any]]:
"""只从卡的历史字段取里程碑,不使用终态摘要字段。"""
fields = payload.get("字段")
if not isinstance(fields, Mapping):
raise AdapterError("卡缺少字段对象,拒绝使用静态卡内容")
for key in ("演变历程", "演变轨迹", "milestones"):
value = fields.get(key)
if isinstance(value, list):
return value
raise AdapterError("卡缺少可按绝对章号冻结的历史字段")
def project_card(row: Mapping[str, Any], *, as_of: int, source_version: str) -> dict[str, Any]:
"""将候选卡投影为截至 as_of 的 eval-only 索引视图。"""
normalized_as_of = _required_chapter(as_of, "as_of")
payload = row.get("draft_payload")
if not isinstance(payload, Mapping):
raise AdapterError(f"卡 {_row_id(row)} 的 draft_payload 不是对象")
card_type = str(payload.get("type") or "")
name = str(payload.get("名称") or "")
if not card_type or not name:
raise AdapterError(f"卡 {_row_id(row)} 缺少 type/名称")
history, omitted = filter_milestones(_card_history(payload), normalized_as_of)
if not history:
raise AdapterError(f"卡 {_row_id(row)} 没有可证明落在 as_of 以前的历史")
appearances: list[int] = []
raw_appearances = payload.get("出场章")
if isinstance(raw_appearances, list):
for value in raw_appearances:
chapter = normalize_chapter(value)
if chapter is not None and chapter <= normalized_as_of:
appearances.append(chapter)
appearances = sorted(set(appearances))
card_id = _row_id(row)
latest = copy.deepcopy(history[-1])
return {
"cardId": card_id,
"type": card_type,
"name": name,
"milestones": copy.deepcopy(history),
"appearanceChapters": appearances,
"derivedState": {
"asOfChapter": normalized_as_of,
"latestMilestone": latest,
},
"source": {
"sourceId": f"eval-draft:{card_id}",
"sourceVersion": source_version,
"scope": "card_projection",
"chapterRange": f"1-{normalized_as_of}",
},
"evaluationStatus": "eval_draft",
"upstreamStatus": str(row.get("status") or "unknown"),
"productionRetrievalEligible": False,
"omittedHistoryCount": len(omitted),
}
def _card_source_version(row: Mapping[str, Any], base_version: str) -> str:
"""把卡行 revision 纳入来源版本,防止卡内容变更复用旧版本。"""
return f"{base_version}:card-{_row_id(row)}-rev-{row.get('revision') or 0}"
def _project_outline(row: Mapping[str, Any]) -> dict[str, Any]:
"""保留窗的结构化摘要和绝对边界,不使用 window_no 作为冻结键。"""
start = _required_chapter(row.get("from_order"), "outline.from_order")
end = _required_chapter(row.get("to_order"), "outline.to_order")
if end < start:
raise AdapterError("大纲窗 from_order 大于 to_order")
return {
"from_order": start,
"to_order": end,
"windowNo": row.get("window_no"),
"outline": str(row.get("outline_text") or ""),
"checkStatus": str(row.get("check_status") or "unknown"),
"sourceId": f"outline-window:{row.get('id')}",
}
def _project_scaffold(row: Mapping[str, Any], source_version: str) -> dict[str, Any]:
"""只取历史章细纲摘要,不查询或复制正文 block。"""
chapter = _required_chapter(row.get("chapter"), "scaffold.chapter")
result: dict[str, Any] = {
"chapter": chapter,
"title": str(row.get("title") or ""),
"outline": str(row.get("outline_text") or ""),
"sourceId": f"scaffold:{row.get('id')}",
"sourceVersion": source_version,
}
if isinstance(row.get("pattern_hints"), list):
result["patternHints"] = copy.deepcopy(row["pattern_hints"])
return result
def _target_facts_from_scaffold(target_scaffold: Mapping[str, Any], target: int) -> dict[str, Any]:
"""从目标章 reference scaffold 生成审计侧 proxy,不送入 planner。"""
target_id = target_scaffold.get("id")
text = str(target_scaffold.get("outline_text") or "").strip()
if not text:
raise AdapterError("目标章 scaffold 缺少结构化事实,不能执行内容级审计")
fragments = [part.strip() for part in re.split(r"[。;;!?!?\n]+", text) if part.strip()]
facts: list[dict[str, Any]] = []
seen: set[str] = set()
for index, fragment in enumerate([text, *fragments]):
if len(fragment) < 4 or fragment in seen:
continue
seen.add(fragment)
facts.append(
{
"id": f"target-scaffold:{target_id}:{index}",
"firstChapter": target,
"text": fragment,
}
)
if not facts:
raise AdapterError("目标章 scaffold 没有可用于审计的结构化事实")
return {
"targetChapter": target,
"source": "reference_scaffold_proxy",
"forbiddenFacts": facts,
}
def _validate_selection(selection: Mapping[str, Any]) -> tuple[list[str], list[str]]:
"""校验正确卡和 placebo 卡的预注册集合。"""
if not isinstance(selection, Mapping):
raise AdapterError("card selection 必须是对象")
correct = _normalize_id_list(selection.get("correctCardIds"), "correctCardIds")
placebo = _normalize_id_list(selection.get("placeboCardIds"), "placeboCardIds")
if set(correct) & set(placebo):
raise AdapterError("correctCardIds 与 placeboCardIds 不能重叠")
return correct, placebo
def build_replay_config(
*,
work: Mapping[str, Any],
reference: Mapping[str, Any],
outline_rows: Sequence[Mapping[str, Any]],
scaffold_rows: Sequence[Mapping[str, Any]],
target_scaffold: Mapping[str, Any],
card_rows: Sequence[Mapping[str, Any]],
card_selection: Mapping[str, Any],
as_of: int,
target: int,
evaluation_set_version: str,
strategy_version: str,
source: Mapping[str, Any],
authorization: Mapping[str, Any],
run_id: str | None = None,
snapshot_version: str = DEFAULT_SNAPSHOT_VERSION,
history_chapter_limit: int = 6,
) -> dict[str, Any]:
"""把只读查询结果组装为 run_replay 可消费的临时配置。"""
normalized_as_of = _required_chapter(as_of, "as_of")
normalized_target = _required_chapter(target, "target")
if normalized_target != normalized_as_of + 1:
raise AdapterError("target 必须等于 as_of+1")
if not str(evaluation_set_version).strip() or not str(strategy_version).strip():
raise AdapterError("evaluation_set_version/strategy_version 不能为空")
target_chapter = _required_chapter(target_scaffold.get("chapter"), "target_scaffold.chapter")
if target_chapter != normalized_target:
raise AdapterError("target scaffold 不是目标章,拒绝混用")
source_hash = str(source.get("sourceHash") or "")
source_version = str(source.get("sourceVersion") or "")
if not source_hash.startswith("sha256:") or source_version != f"raw-file-v1:{source_hash}":
raise AdapterError("来源 hash/version 不符合原文件版本合同")
kept_windows, _ = filter_outline_windows(outline_rows, normalized_as_of)
projected_windows = [_project_outline(row) for row in kept_windows]
historical = []
for row in scaffold_rows:
chapter = _required_chapter(row.get("chapter"), "scaffold.chapter")
if chapter <= normalized_as_of:
historical.append(row)
historical.sort(key=lambda row: _required_chapter(row.get("chapter"), "scaffold.chapter"))
if history_chapter_limit <= 0:
raise AdapterError("history_chapter_limit 必须为正数")
projected_scaffolds = [
_project_scaffold(row, source_version)
for row in historical[-history_chapter_limit:]
]
correct_ids, placebo_ids = _validate_selection(card_selection)
if any(str(row.get("source_type") or "") != "upgrade_book" for row in card_rows):
raise AdapterError("卡选择包含非 upgrade_book 来源")
rows_by_id = {_row_id(row): row for row in card_rows}
if len(rows_by_id) != len(card_rows):
raise AdapterError("卡查询结果含重复 id")
selected_ids = set(correct_ids + placebo_ids)
if set(rows_by_id) != selected_ids:
missing = sorted(selected_ids - set(rows_by_id))
unexpected = sorted(set(rows_by_id) - selected_ids)
raise AdapterError(f"卡查询结果与预注册不一致: missing={missing}, unexpected={unexpected}")
projected_cards = {
card_id: project_card(
rows_by_id[card_id],
as_of=normalized_as_of,
source_version=_card_source_version(rows_by_id[card_id], source_version),
)
for card_id in sorted(selected_ids)
}
correct_cards = [projected_cards[card_id] for card_id in correct_ids]
placebo_cards = [projected_cards[card_id] for card_id in placebo_ids]
source_catalog: list[dict[str, Any]] = [
{
"sourceId": f"reference-work:{work.get('id')}",
"sourceVersion": source_version,
"sourceHash": source_hash,
"scope": "metadata",
"sourceStatus": str(reference.get("parse_status") or "unknown"),
}
]
for row in kept_windows:
source_catalog.append(
{
"sourceId": f"outline-window:{row.get('id')}",
"sourceVersion": source_version,
"from_order": _required_chapter(row.get("from_order"), "outline.from_order"),
"to_order": _required_chapter(row.get("to_order"), "outline.to_order"),
"scope": "outline_window",
}
)
for row in historical[-history_chapter_limit:]:
source_catalog.append(
{
"sourceId": f"scaffold:{row.get('id')}",
"sourceVersion": source_version,
"chapter": _required_chapter(row.get("chapter"), "scaffold.chapter"),
"scope": "chapter",
}
)
for card_id in sorted(selected_ids):
source_catalog.append(
{
"sourceId": f"eval-draft:{card_id}",
"sourceVersion": _card_source_version(rows_by_id[card_id], source_version),
"chapterRange": f"1-{normalized_as_of}",
"scope": "card_projection",
"sourceStatus": "eval_draft",
}
)
recent_source_ids = [item["sourceId"] for item in projected_scaffolds]
common_context = {
"L0": {
"purpose": "offline_evaluation",
"scenario": "fine_outline",
"targetChapter": normalized_target,
"outputContract": "fine_outline_v0",
},
"L1": {
"asOfChapter": normalized_as_of,
"recentScaffoldSourceIds": recent_source_ids,
"historyChapterLimit": history_chapter_limit,
},
"L2": {
"referenceWorkId": str(work.get("id")),
"outlineSourceCount": len(projected_windows),
"sourceVersion": source_version,
"sourceHash": source_hash,
},
"L3": {
"sourceMode": "eval_draft",
"authorizationRequired": True,
"targetChapterAvailableOnlyToAudit": True,
},
}
target_facts = _target_facts_from_scaffold(target_scaffold, normalized_target)
reference_work = {
"id": str(work.get("id")),
"title": str(work.get("title") or ""),
"version": source_version,
"chapterCount": reference.get("imported_chapter_count"),
"declaredChapterCount": reference.get("declared_chapter_count"),
}
run_id = run_id or f"replay-work-{work.get('id')}-{normalized_target}"
return {
"runId": run_id,
"referenceWork": reference_work,
"evaluationSetVersion": str(evaluation_set_version),
"strategyVersion": str(strategy_version),
"targetChapter": normalized_target,
"snapshot": {
"asOfChapter": normalized_as_of,
"snapshotVersion": snapshot_version,
"data": {
"outlineWindows": projected_windows,
"chapters": projected_scaffolds,
"cards": [],
},
},
"sources": source_catalog,
"commonContext": common_context,
"arms": {
"outline_only": {"cards": [], "cardSourceIds": [], "cardStrategy": "none"},
"outline_plus_cards": {
"cards": correct_cards,
"cardSourceIds": [f"eval-draft:{card_id}" for card_id in correct_ids],
"cardStrategy": "correct",
},
"outline_plus_placebo_cards": {
"cards": placebo_cards,
"cardSourceIds": [f"eval-draft:{card_id}" for card_id in placebo_ids],
"cardStrategy": "placebo",
},
},
"authorization": copy.deepcopy(dict(authorization)),
"runPermissions": {
"purpose": "offline_evaluation",
"mode": "dry_run",
"sourceMode": "eval_draft",
"writesFormalData": False,
},
"leakageAudit": {
"auditVersion": "content_fact_audit_v0",
"targetFacts": target_facts,
},
}
def load_reference_rows(
*,
dsn: str,
tenant_id: int,
work_id: int,
as_of: int,
target: int,
card_selection: Mapping[str, Any],
) -> dict[str, Any]:
"""在只读事务中读取组装所需的作品、摘要、卡和目标 proxy。"""
normalized_as_of = _required_chapter(as_of, "as_of")
normalized_target = _required_chapter(target, "target")
correct_ids, placebo_ids = _validate_selection(card_selection)
selected_ids = [int(item) if str(item).isdigit() else item for item in correct_ids + placebo_ids]
with psycopg.connect(dsn, row_factory=dict_row) as conn:
conn.execute("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY")
work = conn.execute(
"""
SELECT id,title,revision,chapter_count,parse_status,import_status
FROM muse_content_work
WHERE tenant_id=%s AND id=%s AND deleted=FALSE
""",
(tenant_id, work_id),
).fetchone()
reference_rows = conn.execute(
"""
SELECT id,work_id,declared_chapter_count,imported_chapter_count,
parse_scope,parse_status,source_file,notes,update_time,deleted
FROM example_reference_work
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE
ORDER BY id
""",
(tenant_id, work_id),
).fetchall()
if work is None:
raise AdapterError("作品或参考作品登记不存在")
reference = _unique_record(reference_rows, "reference.source_file")
import_task_rows = conn.execute(
"""
SELECT id,status,command_id,source_snapshot,deleted
FROM muse_content_import_task
WHERE tenant_id=%s AND work_id=%s AND status='succeeded' AND deleted=FALSE
AND source_snapshot->>'file' = %s
ORDER BY id
""",
(tenant_id, work_id, reference.get("source_file")),
).fetchall()
document_rows = conn.execute(
"""
SELECT id,file_name,file_hash,deleted
FROM muse_knowledge_document
WHERE tenant_id=%s AND file_name=%s AND deleted=FALSE
ORDER BY id
""",
(tenant_id, reference.get("source_file")),
).fetchall()
source = validate_source_records(reference_rows, import_task_rows, document_rows)
authorization_row = conn.execute(
"""
SELECT id,snapshot_version,source_hash,source_version,copyright_status,source_status,
allowed_purpose,forbidden_purpose,authorization_basis,authorized_by,
display_summary,checked_at,expires_at,revalidation_at
FROM example_reference_authorization_snapshot
WHERE tenant_id=%s AND work_id=%s AND source_version=%s
ORDER BY checked_at DESC,id DESC
LIMIT 1
""",
(tenant_id, work_id, source["sourceVersion"]),
).fetchone()
authorization = project_authorization(authorization_row, source)
if normalized_target > int(work.get("chapter_count") or 0) + 1:
raise AdapterError("目标章超出作品导入范围")
outline_rows = conn.execute(
"""
SELECT id,window_no,from_order,to_order,outline_text,check_status
FROM example_parse_outline
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE AND from_order<=%s
ORDER BY from_order,to_order,id
""",
(tenant_id, work_id, normalized_as_of),
).fetchall()
scaffold_rows = conn.execute(
"""
SELECT s.id,s.chapter_id,ch.order_no AS chapter,ch.title,s.outline_text,
s.entities,s.pattern_hints
FROM example_parse_scaffold s
JOIN muse_content_chapter ch ON ch.id=s.chapter_id
WHERE s.tenant_id=%s AND s.work_id=%s AND s.deleted=FALSE
AND ch.tenant_id=%s AND ch.deleted=FALSE AND ch.order_no<=%s
ORDER BY ch.order_no,s.id
""",
(tenant_id, work_id, tenant_id, normalized_as_of),
).fetchall()
target_scaffold = conn.execute(
"""
SELECT s.id,s.chapter_id,ch.order_no AS chapter,ch.title,s.outline_text,
s.entities,s.pattern_hints
FROM example_parse_scaffold s
JOIN muse_content_chapter ch ON ch.id=s.chapter_id
WHERE s.tenant_id=%s AND s.work_id=%s AND s.deleted=FALSE
AND ch.tenant_id=%s AND ch.deleted=FALSE AND ch.order_no=%s
ORDER BY s.id LIMIT 1
""",
(tenant_id, work_id, tenant_id, normalized_target),
).fetchone()
if target_scaffold is None:
raise AdapterError("目标章没有 reference scaffold proxy")
card_rows = conn.execute(
"""
SELECT id,status,source_type,source_id,revision,draft_payload
FROM muse_knowledge_draft
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE
AND source_type='upgrade_book' AND id=ANY(%s)
ORDER BY id
""",
(tenant_id, work_id, selected_ids),
).fetchall()
return {
"work": work,
"reference": reference,
"source": source,
"authorization": authorization,
"outline_rows": outline_rows,
"scaffold_rows": scaffold_rows,
"target_scaffold": target_scaffold,
"card_rows": card_rows,
}
def _parse_args() -> argparse.Namespace:
"""解析只读适配器命令行参数。"""
parser = argparse.ArgumentParser(description="从实验库只读组装回放配置")
parser.add_argument("--dsn", default=DSN)
parser.add_argument("--tenant-id", type=int, default=TENANT_ID)
parser.add_argument("--work-id", type=int, required=True)
parser.add_argument("--as-of", type=int, required=True, dest="as_of")
parser.add_argument("--target-chapter", type=int, required=True, dest="target")
parser.add_argument("--card-selection", type=Path, required=True)
parser.add_argument("--output-dir", type=Path, required=True)
parser.add_argument("--evaluation-set-version", default="deep-space-v0")
parser.add_argument("--strategy-version", default="card-index-outline-v0")
parser.add_argument("--run-id")
parser.add_argument("--history-chapter-limit", type=int, default=6)
return parser.parse_args()
def main() -> int:
"""执行只读查询并把配置写到仓库外临时目录。"""
args = _parse_args()
output_dir = args.output_dir.resolve()
if output_dir.is_relative_to(REPO_ROOT.resolve()):
raise AdapterError("适配器输出目录必须位于仓库外")
selection = json.loads(args.card_selection.read_text(encoding="utf-8"))
rows = load_reference_rows(
dsn=args.dsn,
tenant_id=args.tenant_id,
work_id=args.work_id,
as_of=args.as_of,
target=args.target,
card_selection=selection,
)
config = build_replay_config(
**rows,
card_selection=selection,
as_of=args.as_of,
target=args.target,
evaluation_set_version=args.evaluation_set_version,
strategy_version=args.strategy_version,
run_id=args.run_id,
history_chapter_limit=args.history_chapter_limit,
)
output_dir.mkdir(parents=True, exist_ok=True)
(output_dir / "config.json").write_text(_safe_json(config) + "\n", encoding="utf-8")
(output_dir / "target_proxy.json").write_text(
_safe_json(rows["target_scaffold"]) + "\n", encoding="utf-8"
)
summary = {
"runId": config["runId"],
"workId": args.work_id,
"asOfChapter": args.as_of,
"targetChapter": args.target,
"outlineWindowCount": len(config["snapshot"]["data"]["outlineWindows"]),
"historyChapterCount": len(config["snapshot"]["data"]["chapters"]),
"correctCardCount": len(config["arms"]["outline_plus_cards"]["cards"]),
"placeboCardCount": len(config["arms"]["outline_plus_placebo_cards"]["cards"]),
"cardSourceMode": "eval_draft",
"authorizationStatus": config["authorization"]["sourceStatus"],
"configPath": str(output_dir / "config.json"),
}
(output_dir / "adapter_summary.json").write_text(_safe_json(summary) + "\n", encoding="utf-8")
print(json.dumps(summary, ensure_ascii=False, sort_keys=True))
return 0
if __name__ == "__main__":
try:
raise SystemExit(main())
except (AdapterError, OSError, psycopg.Error) as error:
print(json.dumps({"status": "blocked_adapter", "error": str(error)}, ensure_ascii=False))
raise SystemExit(2)