""" 清洗工作流测试(三层架构版本) 测试策略:数据驱动 + 隔离临时任务,覆盖: 1. 底层单词分析(TangoAnalyser):词条 → 状态分类 2. 中间层文件处理(TaskProcessor):文件搬运、去重、格式校验 3. 工作流编排(CleanerWorkflow):状态机推进、人工介入门禁 所有测试用隔离临时目录(tmp_path),不碰真实数据。 """ from pathlib import Path import pytest from pl_japanese.cleaner import ( TangoAnalyser, AnalysisStatus, TaskProcessor, TaskManager, CleanerWorkflow, ACTION_PROCESSED, ACTION_NEED_HUMAN, ACTION_COMPLETED, ACTION_EMPTY, ) from pl_japanese.cleaner.task import ( STATUS_CREATED, STATUS_PROCESSING, STATUS_REVIEWING, STATUS_READY, STATUS_MERGED, ) # ========================================================================== # 测试辅助 # ========================================================================== def _make_task(tmp_path: Path, source_lines: list[str], task_id="test_task"): """创建隔离临时任务(数据源、任务目录、main/skipped 都在 tmp_path 下)""" source = tmp_path / "source.txt" source.write_text("\n".join(source_lines) + "\n", encoding="utf-8", newline="\n") tm = TaskManager(tasks_root=str(tmp_path / "tasks")) task = tm.create_task( task_id=task_id, source=str(source), name=f"test_{task_id}", start_line=1, count=len(source_lines), main=str(tmp_path / "main.txt"), skipped=str(tmp_path / "skipped.txt"), ) return tm, task # ========================================================================== # 1. 底层单词分析(TangoAnalyser):词条 → 状态分类 # ========================================================================== # (kanji, kana, 期望状态, 输出子串检查) ANALYSE_CASES = [ pytest.param("日本人", "にほんじん", AnalysisStatus.SUCCESS, "日|本|人:", id="success-basic"), pytest.param("中国人", "ちゅうごくじん", AnalysisStatus.SUCCESS, "中|国|人:", id="success-polyphone-resolved"), pytest.param("IT", "アイティー", AnalysisStatus.SKIP, "IT:", id="skip-no-kanji"), pytest.param("行列", "こうれつ", AnalysisStatus.SUCCESS, "行|列:こう|れつ:hang|lie", id="polyphone-resolved-by-dict"), pytest.param("女将", "おかみ", AnalysisStatus.SPLIT_FAILED, "女将:", id="split-failed"), pytest.param("食べる", "たべる", AnalysisStatus.SUCCESS, "食|べる:た|べる:", id="verb-auto"), pytest.param("Uターン", "ユーターン", AnalysisStatus.SKIP, "Uターン:", id="skip-no-kanji-with-latin"), ] @pytest.mark.parametrize("kanji,kana,status,substring", ANALYSE_CASES) def test_analyser_status_classification(kanji, kana, status, substring): """底层分析器:词条 → 状态分类正确,输出格式符合预期""" analyser = TangoAnalyser() result = analyser.analyze(kanji, kana) assert result.status == status assert substring in result.formatted_line def test_analyser_needs_review_logic(): """needs_review 属性:SUCCESS/SKIP 返回 False,其余返回 True""" analyser = TangoAnalyser() success = analyser.analyze("日本人", "にほんじん") assert success.is_success is True assert success.needs_review is False skip = analyser.analyze("IT", "アイティー") assert skip.status == AnalysisStatus.SKIP assert skip.needs_review is False # 用一个真正无法自动处理的词(分割失败) split_fail = analyser.analyze("女将", "おかみ") assert split_fail.needs_review is True # ========================================================================== # 2. 中间层文件处理(TaskProcessor):process_source 条数铁律 # ========================================================================== # (源数据, 期望有效行数, 期望输出总数, 期望桶计数子集) PROCESS_CASES = [ pytest.param( ["日本人:にほんじん:", "中国人:ちゅうごくじん:"], 2, 2, {"auto": 2, "skip": 0}, id="all-auto", ), pytest.param( ["中国人:ちゅうごくじん:", "", "学生:がくせい:"], 2, 2, {"auto": 2}, id="blank-ignored", ), pytest.param( ["IT:アイティー:", "%:パーセント:"], 2, 2, {"skip": 2, "auto": 0}, id="all-skip", ), pytest.param( ["日本人:にほんじん:", "IT:アイティー:"], 2, 2, {"auto": 1, "skip": 1}, id="mixed", ), pytest.param( ["行列:こうれつ:", "Uターン:ユーターン:"], 2, 2, {"auto": 1, "skip": 1}, id="mixed-auto-skip", ), ] @pytest.mark.parametrize("lines,valid,total,bucket_subset", PROCESS_CASES) def test_processor_count_law(tmp_path, lines, valid, total, bucket_subset): """中间层文件处理器:条数铁律(有效输入 == 输出总数)+ 桶分布""" tm, task = _make_task(tmp_path, lines) proc = TaskProcessor(task) result = proc.process_source() assert result["valid_lines"] == valid assert result["output_total"] == total assert result["validation"] is True for bucket, cnt in bucket_subset.items(): assert result["buckets"][bucket] == cnt def test_processor_append_mode(tmp_path): """process_source 追加模式:多次调用累加到桶文件""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:", "中国人:ちゅうごくじん:"]) proc = TaskProcessor(task) r1 = proc.process_source(count=1) assert r1["output_total"] == 1 assert proc.snapshot_counts()["auto_done"] == 1 r2 = proc.process_source(start_line=2, count=1) assert r2["output_total"] == 1 assert proc.snapshot_counts()["auto_done"] == 2 # ========================================================================== # 3. 中间层:merge_final 去重 + 格式校验 # ========================================================================== def test_merge_deduplication(tmp_path): """merge_final 去重:重复词条只保留一份""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:", "日本人:にほんじん:"]) proc = TaskProcessor(task) proc.process_source() result = proc.merge_final(dry_run=False) assert result["merged_auto"] == 1 # 2 条输入去重后只合并 1 条 assert result["duplicated_auto"] == 1 assert result["main_total"] == 1 def test_merge_idempotent(tmp_path): """幂等性:同一批数据再次合并,主库不增长""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"]) proc = TaskProcessor(task) proc.process_source() r1 = proc.merge_final(dry_run=False) assert r1["merged_auto"] == 1 proc.clear_single_batch() proc.process_source() # 再次处理同样数据 r2 = proc.merge_final(dry_run=False) assert r2["merged_auto"] == 0 # 无新增 assert r2["duplicated_auto"] == 1 # 全部重复 assert r2["main_total"] == 1 # 总数不变 def test_merge_validation_rejects_bad_format(tmp_path): """merge_final 格式校验:非法行被 reject,不进主库""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"]) proc = TaskProcessor(task) # 先不 process,直接手工写 auto_done 来测试校验 # 手工塞一条合法行(段数匹配) + 一条格式非法行(只有 1 个冒号) proc.auto_done.write_text( "日|本|人:に|ほん|じん:ri|ben|ren\n非法行\n", encoding="utf-8", newline="\n", ) result = proc.merge_final(dry_run=False) assert result["rejected"] == 1 assert result["merged_auto"] == 1 # 合法的那条成功 assert any("非法行" in r[2] for r in result["rejects"]) # ========================================================================== # 4. 中间层:process_review 重跑 # ========================================================================== # (塞入 review 桶名, 内容, 期望重跑条数, 期望落到的桶) REPROCESS_CASES = [ pytest.param("pinyin", "中国人:ちゅうごくじん:", 1, "auto", id="pinyin-to-auto"), pytest.param("split", "日本人:にほんじん:", 1, "auto", id="split-to-auto"), pytest.param("special", "IT:アイティー:", 1, "skip", id="special-to-skip"), ] @pytest.mark.parametrize("bucket,content,reprocessed,land_bucket", REPROCESS_CASES) def test_process_review_reclassifies(tmp_path, bucket, content, reprocessed, land_bucket): """process_review:重新分类后清空原桶,条目重新分流""" tm, task = _make_task(tmp_path, ["先生:せんせい:"]) proc = TaskProcessor(task) proc.process_source() # 手工塞一条到 review 桶 proc.review_files[bucket].write_text(content + "\n", encoding="utf-8", newline="\n") result = proc.process_review(bucket) assert result["reprocessed"] == reprocessed assert result["new_distribution"][land_bucket] >= 1 # 原 review 桶已清零 assert len(proc._read_clean_lines(proc.review_files[bucket])) == 0 def test_process_review_unknown_bucket_raises(tmp_path): """process_review 非法桶名抛异常""" tm, task = _make_task(tmp_path, ["先生:せんせい:"]) proc = TaskProcessor(task) with pytest.raises(ValueError, match="未知 review 桶"): proc.process_review("unknown_bucket") # ========================================================================== # 5. 工作流编排(CleanerWorkflow):状态机推进 # ========================================================================== def test_workflow_run_auto_complete(tmp_path): """全自动流程:created → ready(无 review)→ merged → completed""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:", "中国人:ちゅうごくじん:"]) wf = CleanerWorkflow(task) # 第一步:created → process → ready(全部 auto,无 review) r1 = wf.run() assert r1["action"] == ACTION_PROCESSED assert r1["status"] == STATUS_READY assert r1["review_total"] == 0 # 第二步:ready → merge → merged r2 = wf.run() assert r2["action"] == ACTION_PROCESSED assert r2["status"] == STATUS_MERGED assert r2["merge"]["merged_auto"] == 2 # 第三步:merged → completed(终态) r3 = wf.run() assert r3["action"] == ACTION_COMPLETED assert r3["status"] == STATUS_MERGED def test_workflow_run_with_review_gate(tmp_path): """带人工介入流程:created → reviewing → 返回待人工(不推进)""" # 使用真正会进 review 的词:女将(分割失败) tm, task = _make_task(tmp_path, ["女将:おかみ:", "日本人:にほんじん:"]) wf = CleanerWorkflow(task) # 第一步:created → process → reviewing(有 review_split) r1 = wf.run() assert r1["action"] == ACTION_PROCESSED assert r1["status"] == STATUS_REVIEWING assert r1["review_total"] > 0 # 第二步:reviewing 且 review 非空 → 返回待人工,不推进、不写文件 r2 = wf.run() assert r2["action"] == ACTION_NEED_HUMAN assert r2["status"] == STATUS_REVIEWING assert "review" in r2 assert r2["review_total"] > 0 def test_workflow_review_cleared_advances_to_ready(tmp_path): """人工重跑清零 review 后,run() 自动推进到 ready""" # 使用会进 review_split 的词 tm, task = _make_task(tmp_path, ["女将:おかみ:"]) wf = CleanerWorkflow(task) proc = wf.processor # 处理 → reviewing r1 = wf.run() assert r1["status"] == STATUS_REVIEWING # 人工重跑 review_split(假设改了字典,重跑后全部转 auto) proc.review_files["split"].write_text("", encoding="utf-8", newline="\n") proc.auto_done.write_text("女|将:お|かみ:nv|jiang\n", encoding="utf-8", newline="\n") task.state.bucket_counts = proc.snapshot_counts() task.save() # 再次 run():review 清零 → 转 ready(不合并,把合并留给下一次) r2 = wf.run() assert r2["action"] == ACTION_PROCESSED assert r2["status"] == STATUS_READY def test_workflow_merge_rejects_when_review_reopened(tmp_path): """合并前门禁:ready 时 review 又有内容(被人工重新塞入),拒绝合并""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"]) wf = CleanerWorkflow(task) proc = wf.processor # 处理 → ready r1 = wf.run() assert r1["status"] == STATUS_READY # 人工重新塞一条 review(模拟发现新问题) proc.review_files["pinyin"].write_text( "测试:てすと:ce|shi\n", encoding="utf-8", newline="\n", ) # run() 合并被拒绝,状态回到 reviewing r2 = wf.run() assert r2["action"] == ACTION_NEED_HUMAN assert r2["status"] == STATUS_REVIEWING def test_workflow_merge_rejects_bad_format(tmp_path): """合并时格式校验失败:不推进状态,保持 ready,提示人工修正""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"]) wf = CleanerWorkflow(task) proc = wf.processor r1 = wf.run() assert r1["status"] == STATUS_READY # 手工塞一条非法格式 proc.auto_done.write_text( "日|本|人:にほん|じん:ri|ben|ren\n非法行", encoding="utf-8", newline="\n", ) r2 = wf.run() assert r2["action"] == ACTION_NEED_HUMAN assert r2["status"] == STATUS_READY assert "格式非法" in r2["message"] def test_workflow_empty_source(tmp_path): """空源文件:返回 empty,不推进""" tm, task = _make_task(tmp_path, ["", " "]) wf = CleanerWorkflow(task) result = wf.run() assert result["action"] == ACTION_EMPTY assert result["process"]["valid_lines"] == 0 # ========================================================================== # 6. 任务管理(TaskManager):CRUD + 多任务隔离 # ========================================================================== def test_task_manager_create_and_load(tmp_path): """任务创建和加载:配置持久化正确""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"], task_id="task1") loaded = tm.load_task("task1") assert loaded.task_id == "task1" assert loaded.config.source == task.config.source assert loaded.state.status == STATUS_CREATED def test_task_manager_list_summaries(tmp_path): """list_task_summaries:返回所有任务摘要""" source1 = tmp_path / "s1.txt" source1.write_text("日本人:にほんじん:\n", encoding="utf-8") source2 = tmp_path / "s2.txt" source2.write_text("中国人:ちゅうごくじん:\n", encoding="utf-8") tm = TaskManager(tasks_root=str(tmp_path / "tasks")) tm.create_task("t1", source=str(source1), name="任务1") tm.create_task("t2", source=str(source2), name="任务2") summaries = tm.list_task_summaries() assert len(summaries) == 2 ids = {s["task_id"] for s in summaries} assert ids == {"t1", "t2"} def test_task_manager_duplicate_rejected(tmp_path): """重复 task_id 创建被拒绝""" tm, task = _make_task(tmp_path, ["日本人:にほんじん:"]) with pytest.raises(FileExistsError): tm.create_task(task.task_id, source=task.config.source) def test_multi_task_isolation(tmp_path): """多任务隔离:各任务的单批文件独立,互不干扰""" source = tmp_path / "source.txt" source.write_text("日本人:にほんじん:\n中国人:ちゅうごくじん:\n", encoding="utf-8") tm = TaskManager(tasks_root=str(tmp_path / "tasks")) t1 = tm.create_task("t1", source=str(source), count=1, main=str(tmp_path / "main.txt"), skipped=str(tmp_path / "skip.txt")) t2 = tm.create_task("t2", source=str(source), start_line=2, count=1, main=str(tmp_path / "main.txt"), skipped=str(tmp_path / "skip.txt")) p1 = TaskProcessor(t1) p2 = TaskProcessor(t2) p1.process_source() p2.process_source() # 各自单批文件独立 assert p1.snapshot_counts()["auto_done"] == 1 assert p2.snapshot_counts()["auto_done"] == 1 # 但 main 库是共享的(都指向同一文件) p1.merge_final(dry_run=False) p2.merge_final(dry_run=False) assert p1.final_counts()["main"] == 2 assert p2.final_counts()["main"] == 2