diff --git a/app/question_gen/pair_atomic_writer.py b/app/question_gen/pair_atomic_writer.py index 8660da8..7b6412c 100644 --- a/app/question_gen/pair_atomic_writer.py +++ b/app/question_gen/pair_atomic_writer.py @@ -59,7 +59,7 @@ class PairPendingBuffer: """ def __init__(self) -> None: - """初始化空缓冲器。pending 以 pair_id 索引首个到达的孪生对成员。""" + """pending 以 pair_id 索引首个到达的孪生对成员,等伙伴到齐再 emit。""" self._pending: dict[str, GeneratedQuestion] = {} def add(self, q: GeneratedQuestion) -> QuestionUnit | None: @@ -89,11 +89,7 @@ class PairPendingBuffer: return unit def pending_orphans(self) -> list[GeneratedQuestion]: - """返回当前仍未配齐的孤儿成员(供批次末尾检测悬挂项)。 - - 返回: - 尚在缓冲、未等到伙伴的孪生对成员列表(按 pair_id 插入顺序)。 - """ + """返回仍未配齐的悬挂成员,供调用方在批次末尾检测"只落 P 未落 Q"。""" return list(self._pending.values()) @@ -138,7 +134,7 @@ def _validate_disk_pair(pair_id: str, group: list[GeneratedQuestion]) -> list[Ge ) -def _sift_disk_pairs(questions: list[GeneratedQuestion]) -> list[GeneratedQuestion]: +def _keep_complete_disk_pairs(questions: list[GeneratedQuestion]) -> list[GeneratedQuestion]: """筛选磁盘读回的题目:single 全保留、pair 按 ``_validate_disk_pair`` 三态处理。 悬挂孤儿 warn+drop、结构损坏/绑定不一致 raise ValueError、合法成对保留后交给 @@ -182,7 +178,7 @@ def write_accepted(path: Path, units: list[QuestionUnit]) -> None: def read_accepted(path: Path) -> list[QuestionUnit]: """读回 path 的 accepted JSON → 聚合为单元列表。 - 磁盘是外部输入,按 P5 全量校验后再用:``_sift_disk_pairs`` 对 pair 分组三态处理 + 磁盘是外部输入,按 P5 全量校验后再用:``_keep_complete_disk_pairs`` 对 pair 分组三态处理 ——"只落 P 未落 Q"的悬挂孤儿 warn+drop(不进结果、不 raise);结构损坏(超员/ 角色重复)或绑定不一致(video_id/task_type/flip_axis)显式 raise ValueError(不 依赖 build_units 内会被 ``-O`` 剥除的 assert)。single 全保留。 @@ -196,5 +192,5 @@ def read_accepted(path: Path) -> list[QuestionUnit]: raw = json.loads(path.read_text(encoding="utf-8")) questions = [_dict_to_q(d) for d in raw] - kept = _sift_disk_pairs(questions) + kept = _keep_complete_disk_pairs(questions) return validate_units(build_units(kept)) diff --git a/tests/unit/test_pair_atomic_write.py b/tests/unit/test_pair_atomic_write.py index 71c9c6c..9113552 100644 --- a/tests/unit/test_pair_atomic_write.py +++ b/tests/unit/test_pair_atomic_write.py @@ -15,6 +15,7 @@ from __future__ import annotations import json import pytest +from loguru import logger from app.harness.pools import _q_to_dict from app.question_gen.pair_atomic_writer import ( @@ -179,7 +180,7 @@ def test_write_read_roundtrip_keeps_pair_and_single(tmp_path) -> None: assert {q.flip_axis for q in loaded_pair.questions} == {"before_after"} -def test_read_drops_disk_dangling_orphan(tmp_path, caplog) -> None: +def test_read_drops_disk_dangling_orphan(tmp_path) -> None: """磁盘上某 pair 只有 original(缺 mirror)→ read_accepted warn+drop,不 raise。""" single = _make_single() original, _ = _make_pair() @@ -188,11 +189,19 @@ def test_read_drops_disk_dangling_orphan(tmp_path, caplog) -> None: out = tmp_path / "bank.json" out.write_text(json.dumps(records, ensure_ascii=False, indent=2), encoding="utf-8") - loaded = read_accepted(out) + # loguru 不走标准 logging,用项目既定 sink 捕获模式(见 test_inference_pair_aggregate) + captured: list[str] = [] + sink_id = logger.add(captured.append, level="WARNING", format="{message}") + try: + loaded = read_accepted(out) + finally: + logger.remove(sink_id) assert len(loaded) == 1 assert loaded[0].kind == "single" assert loaded[0].questions[0].question_id == single.question_id + # warn+drop 契约的 "warn" 半:悬挂孤儿必须告警,不静默丢弃 + assert any("悬挂孤儿" in msg for msg in captured), "悬挂孤儿未告警(静默丢弃)" def test_read_raises_on_disk_corrupt_duplicate_role(tmp_path) -> None: