feat: add flip gate reusing P prediction and mirror-Q agent run

This commit is contained in:
2026-07-14 16:26:51 -04:00
parent 8731e448fe
commit 15aee0cfc1
2 changed files with 504 additions and 15 deletions
+262 -15
View File
@@ -13,7 +13,7 @@ import enum
import hashlib
import json
from pathlib import Path
from typing import TYPE_CHECKING, Protocol
from typing import TYPE_CHECKING, Protocol, overload
from json_repair import repair_json
from loguru import logger
@@ -38,17 +38,39 @@ class FlipDecision(enum.Enum):
FLIP_SKIPPED = "flip_skipped"
def question_hash(question: str, options: tuple[str, ...], answer: str) -> str:
@overload
def question_hash(question: GeneratedQuestion) -> str: ...
@overload
def question_hash(question: str, options: tuple[str, ...], answer: str) -> str: ...
def question_hash(
question: GeneratedQuestion | str,
options: tuple[str, ...] | None = None,
answer: str | None = None,
) -> str:
"""题 payloadquestion+options+answer)的稳定 hash,防 JSON 变动误用旧 verdict。
两种等价调用形态:整题 `question_hash(q)` 或散参 `question_hash(题面, 选项, 答案)`
前者按题面/选项/答案拆解后走同一路径,保证与散参形态哈希一致。
参数:
question: 题目文本。
options: 选项元组。
answer: 正确答案字母。
question: 整条 GeneratedQuestion,或题目文本字符串
options: 选项元组(散参形态必传)
answer: 正确答案字母(散参形态必传)
返回:
16 位十六进制摘要。
异常:
TypeError: 传入题面字符串却缺 options / answer(散参形态参数不全)。
"""
if isinstance(question, GeneratedQuestion):
return question_hash(question.question, question.options, question.answer)
if options is None or answer is None:
raise TypeError("散参形态 question_hash 需同时传入 (question, options, answer)")
payload = json.dumps(
{
"question": question,
@@ -90,7 +112,7 @@ def canonical_answer_text(options: tuple[str, ...], letter: str | None) -> str |
return None
opt = options[idx]
prefix = f"{s}. "
return opt[len(prefix):] if opt.startswith(prefix) else opt
return opt[len(prefix) :] if opt.startswith(prefix) else opt
def judge_flip(*, p_text: str | None, q_text: str | None) -> FlipDecision:
@@ -131,8 +153,8 @@ class AgentRunner(Protocol):
def _cheat_hash(question: GeneratedQuestion) -> str:
"""按散参数签名计算作弊门题面 hashquestion/options/answer)。"""
return question_hash(question.question, question.options, question.answer)
"""计算作弊门题面 hashquestion/options/answer,供两门与终判统一续跑主键"""
return question_hash(question)
def _recover_survivors(
@@ -208,15 +230,24 @@ async def run_cheater_gate(
correct = pred is not None and pred.strip().upper() == q.answer.strip().upper()
verdict = "filtered_too_easy" if correct else "passed"
store.record_verdict(
question_id=q.question_id, question_hash=_cheat_hash(q), stage="cheat",
round=round_no, agent_prediction=pred, agent_correct=correct,
verdict=verdict, pair_id=None, agent_config=cfg_fp,
question_id=q.question_id,
question_hash=_cheat_hash(q),
stage="cheat",
round=round_no,
agent_prediction=pred,
agent_correct=correct,
verdict=verdict,
pair_id=None,
agent_config=cfg_fp,
)
if not correct:
survivors.append(q)
logger.info(
"作弊门: {} 题(续跑复用 {},新判 {})→ 存活 {}",
len(questions), len(completed), len(todo), len(survivors),
len(questions),
len(completed),
len(todo),
len(survivors),
)
return survivors
@@ -293,9 +324,7 @@ def _existing_frames(frame_paths: list[str]) -> list[str]:
return [p for p in frame_paths if Path(p).exists()]
def _build_mirror_question(
question: GeneratedQuestion, mirror: dict
) -> GeneratedQuestion | None:
def _build_mirror_question(question: GeneratedQuestion, mirror: dict) -> GeneratedQuestion | None:
"""从解析出的 mirror dict 构造镜像题;字段缺失/类型错误 → None。"""
try:
options = tuple(str(o) for o in mirror["options"])
@@ -371,3 +400,221 @@ async def generate_mirror_question(
)
return None
return mirror_q
def _read_cheat_prediction(
store: QuestionGenStore, q: GeneratedQuestion, cfg_fp: str
) -> str | None:
"""从 adversarial_verdicts 读作弊门落库的 P 预测字母(C2:绝不重跑 agent)。
按 (question_id, 当前题面 hash, stage='cheat', 当前 agent_config) 定位那条
由 `run_cheater_gate` 写入的预测;无匹配行返回 None。
"""
row = store._conn.execute(
"SELECT agent_prediction FROM adversarial_verdicts "
"WHERE question_id=? AND question_hash=? AND stage='cheat' AND agent_config=?",
(q.question_id, _cheat_hash(q), cfg_fp),
).fetchone()
return row[0] if row else None
def _persist_flip(
store: QuestionGenStore,
q: GeneratedQuestion,
decision: FlipDecision,
mirror_pred: str | None,
round_no: int,
cfg_fp: str,
pair_id: str,
) -> None:
"""写 flip_original + flip_mirror 两条 verdict,并按判定改写原题 cheat 行。
flip_original 复用 P 的作弊门预测;flip_mirror 记镜像预测;两者同 pair_id 关联。
FILTERED_NO_FLIP 时把 cheat 行 verdict 改判 filtered_no_flip(与终判排除双保险);
passed / flip_skipped 时 cheat 行保持 passed。
"""
h = _cheat_hash(q)
p_pred = _read_cheat_prediction(store, q, cfg_fp)
store.record_verdict(
question_id=q.question_id,
question_hash=h,
stage="flip_original",
round=round_no,
agent_prediction=p_pred,
agent_correct=None,
verdict=decision.value,
pair_id=pair_id,
agent_config=cfg_fp,
)
store.record_verdict(
question_id=q.question_id,
question_hash=h,
stage="flip_mirror",
round=round_no,
agent_prediction=mirror_pred,
agent_correct=None,
verdict=decision.value,
pair_id=pair_id,
agent_config=cfg_fp,
)
if decision is FlipDecision.FILTERED_NO_FLIP:
store.record_verdict(
question_id=q.question_id,
question_hash=h,
stage="cheat",
round=round_no,
agent_prediction=p_pred,
agent_correct=False,
verdict="filtered_no_flip",
pair_id=None,
agent_config=cfg_fp,
)
async def _judge_one_flip(
q: GeneratedQuestion,
flip_axis: str | None,
*,
agent: AgentRunner,
vlm: VLMProvider,
trees: dict[str, TreeIndex],
store: QuestionGenStore,
cfg_fp: str,
config: AdversarialFilterConfig,
run_id: str,
session_id: str,
) -> tuple[FlipDecision, str | None]:
"""跑单题翻转判定,返回 (decision, 镜像预测字母)。
原题 P 预测**只从 adversarial_verdicts 表读作弊门落的行**(不重跑 agent),
故 `store` 与 `cfg_fp` 必传(C2:按 (question_id, question_hash, stage='cheat',
agent_config) 定位那条预测)。镜像造不出 / 素材缺失 → FLIP_SKIPPED(不跑 agent)。
"""
tree = trees.get(q.video_id)
if tree is None or flip_axis is None:
return FlipDecision.FLIP_SKIPPED, None
material = _rebuild_material(tree, q.source_nodes)
mirror = await generate_mirror_question(
q, flip_axis=flip_axis, vlm=vlm, material=material, session_id=session_id
)
if mirror is None:
return FlipDecision.FLIP_SKIPPED, None
preds = await agent.predict(
[mirror], max_steps=config.adversarial_agent_max_steps, run_id=f"{run_id}_mirror"
)
q_pred = preds.get(mirror.question_id)
p_pred = _read_cheat_prediction(store, q, cfg_fp) # 复用作弊门 P 预测(不重跑)
p_text = canonical_answer_text(q.options, p_pred)
q_text = canonical_answer_text(mirror.options, q_pred)
return judge_flip(p_text=p_text, q_text=q_text), q_pred
async def run_flip_gate(
survivors: list[GeneratedQuestion],
*,
agent: AgentRunner,
vlm: VLMProvider,
store: QuestionGenStore,
trees: dict[str, TreeIndex],
config: AdversarialFilterConfig,
round_no: int,
run_id: str,
session_id: str,
) -> list[GeneratedQuestion]:
"""翻转门:不支持 flip 的终判 passed;支持的按 canonical 翻转判定。
P 预测复用作弊门落表结果(不重跑);仅新跑镜像 Q。任一无效 / 镜像失败 →
flip_skipped(保留题,只经作弊门,不误杀)。镜像题只用于判定,不进题库。
参数:
survivors: 作弊门存活(agent 答错)的题列表。
agent: 完整 agent 试答端口(仅对镜像 Q 调用)。
vlm: 镜像题生成 VLM 端口。
store: verdict 持久化(含作弊门 P 预测来源)。
trees: video_id → 三层树索引(重建镜像素材用)。
config: 过滤配置(提供 max_steps)。
round_no: 当前轮次(构造 pair_id)。
run_id: agent 推理 run 标识。
session_id: VLM 遥测会话 ID。
返回:
终判 verdict∈{passed, flip_skipped} 的题(filtered_no_flip 被剔除)。
"""
from app.question_gen.strategy_action_recognition import _AR_PATTERN_BY_NAME
cfg_fp = agent_config_fingerprint(
skill_mode=agent.skill_mode,
max_steps=config.adversarial_agent_max_steps,
model=agent.model,
)
kept: list[GeneratedQuestion] = []
for q in survivors:
sp = _AR_PATTERN_BY_NAME.get(q.sub_pattern or "")
if sp is None or not sp.supports_flip:
kept.append(q) # cheat 已记 passed,无需改写
continue
decision, mirror_pred = await _judge_one_flip(
q,
sp.flip_axis,
agent=agent,
vlm=vlm,
trees=trees,
store=store,
cfg_fp=cfg_fp,
config=config,
run_id=run_id,
session_id=session_id,
)
pair_id = f"{q.question_id}::{round_no}"
_persist_flip(store, q, decision, mirror_pred, round_no, cfg_fp, pair_id)
if decision is not FlipDecision.FILTERED_NO_FLIP:
kept.append(q) # passed 或 flip_skipped 都保留
logger.info("翻转门: {} 存活 → 保留 {}", len(survivors), len(kept))
return kept
def _question_to_record(q: GeneratedQuestion) -> dict:
"""把题目序列化为最终题库 JSON 记录(字段与 loader.load_benchmark 读回口径一致)。"""
return {
"question_id": q.question_id,
"video_id": q.video_id,
"task_type": q.task_type,
"question": q.question,
"options": list(q.options),
"answer": q.answer,
"source_nodes": list(q.source_nodes),
"difficulty": q.difficulty,
"family": q.family,
"skill_target": q.skill_target,
"difficulty_steps": q.difficulty_steps,
"sub_pattern": q.sub_pattern,
}
def write_final_bank(
out_path: Path,
store: QuestionGenStore,
questions_by_id: dict[str, GeneratedQuestion],
cfg_fp: str,
) -> list[dict]:
"""按终判 passed 全量重建最终题库 JSON(镜像题绝不入库)。
终判集合来自 `final_passed_question_ids`(当前 hash+config 下 cheat=passed 且
无 filtered_no_flip 行)。镜像题不在 questions_by_id 中,天然被排除。
参数:
out_path: 输出 JSON 路径。
store: verdict 来源。
questions_by_id: question_id → 原题(仅 Phase A 产物,不含镜像)。
cfg_fp: 当前 agent 配置指纹。
返回:
写入的记录列表(保持 questions_by_id 的插入顺序)。
"""
hash_by_qid = {qid: _cheat_hash(q) for qid, q in questions_by_id.items()}
passed = store.final_passed_question_ids(hash_by_qid, cfg_fp)
records = [
_question_to_record(questions_by_id[qid]) for qid in questions_by_id if qid in passed
]
out_path.write_text(json.dumps(records, ensure_ascii=False, indent=2), encoding="utf-8")
return records