""" 清洗工作流测试(任务驱动 + 人工介入协作) 本文件把"处理 → review → 重跑 → 合并"协作流程写成可运行的测试。 每个测试用隔离的临时任务目录(tmp_path),不碰真实数据(legacy_batch1 / data/db)。 设计要点(体现工作流本质): - 工作流需要人工介入:review 桶里的条目由人核对,测试用"模拟人工确认"占位 - merge 有授权门禁:review 未清零时 merge_all 必须拒绝 - 条数铁律:输入有效行数 == 各桶输出之和 - 状态机流转:created → processing → reviewing → ready → merged 日常协作(真实数据)请直接用 driver 函数在这里加临时测试驱动, 或用 `python -m pl_japanese.cleaner.cli` / `jclean` 命令行。 """ from pathlib import Path import pytest from pl_japanese.cleaner import BatchProcessor, TaskManager, Task from pl_japanese.cleaner.task import ( STATUS_CREATED, STATUS_REVIEWING, STATUS_READY, STATUS_MERGED, ) # -------------------------------------------------------------------------- # 夹具:隔离的任务环境 # -------------------------------------------------------------------------- @pytest.fixture def isolated_env(tmp_path): """ 构建一个完全隔离的任务环境: - tasks_root 在 tmp_path 下 - 权威库(main/skipped)也在 tmp_path 下(不碰 data/db) - 数据源是临时写入的小样本 返回 (tasks_root, main, skipped, make_source) """ tasks_root = tmp_path / "tasks" main = tmp_path / "db" / "vocabulary.txt" skipped = tmp_path / "db" / "skipped.txt" main.parent.mkdir(parents=True, exist_ok=True) main.write_text("", encoding="utf-8", newline="\n") skipped.write_text("", encoding="utf-8", newline="\n") def make_source(lines, name="source.txt"): p = tmp_path / name p.write_text("\n".join(lines) + "\n", encoding="utf-8", newline="\n") return p return tasks_root, main, skipped, make_source def _new_task(isolated_env, source_lines, count=100): tasks_root, main, skipped, make_source = isolated_env src = make_source(source_lines) tm = TaskManager(tasks_root=str(tasks_root)) task = tm.create_task( task_id="t1", source=str(src), name="测试任务", start_line=1, count=count, main=str(main), skipped=str(skipped), ) return tm, task # -------------------------------------------------------------------------- # 任务生命周期 # -------------------------------------------------------------------------- def test_task_create_and_load(isolated_env): """创建任务后可从磁盘重新加载,配置一致""" tm, task = _new_task(isolated_env, ["日本語:にほんご:"]) assert task.state.status == STATUS_CREATED reloaded = tm.load_task("t1") assert reloaded.task_id == "t1" assert reloaded.config.source == task.config.source assert reloaded.config.main == task.config.main assert reloaded.state.status == STATUS_CREATED def test_multi_task_isolation(isolated_env): """多任务并存:各自独立目录,互不干扰""" tasks_root, main, skipped, make_source = isolated_env tm = TaskManager(tasks_root=str(tasks_root)) s1 = make_source(["学生:がくせい:"], "s1.txt") s2 = make_source(["先生:せんせい:"], "s2.txt") t1 = tm.create_task("task_a", str(s1), main=str(main), skipped=str(skipped)) t2 = tm.create_task("task_b", str(s2), main=str(main), skipped=str(skipped)) assert Path(t1.task_dir) != Path(t2.task_dir) ids = {s["task_id"] for s in tm.list_task_summaries()} assert ids == {"task_a", "task_b"} def test_create_duplicate_rejected(isolated_env): """重复 task_id 创建被拒绝""" tm, task = _new_task(isolated_env, ["日本:にほん:"]) with pytest.raises(FileExistsError): tm.create_task("t1", source=task.config.source) # -------------------------------------------------------------------------- # 处理与状态流转 # -------------------------------------------------------------------------- def test_process_count_invariant(isolated_env): """条数铁律:有效输入行数 == 各桶输出之和""" lines = [ "中国人:ちゅうごくじん:", "日本人:にほんじん:", "学生:がくせい:", "", # 空行不计 "先生:せんせい:", ] tm, task = _new_task(isolated_env, lines) bp = BatchProcessor(task) result = bp.process_batch() assert result["valid_lines"] == 4 assert result["output_total"] == 4 assert result["validation"] is True def test_process_all_auto_goes_ready(isolated_env): """全部自动处理(无 review)→ 任务状态直接 ready""" lines = ["日本人:にほんじん:", "中国人:ちゅうごくじん:"] tm, task = _new_task(isolated_env, lines) bp = BatchProcessor(task) result = bp.process_batch() assert result["buckets"]["auto"] == 2 assert result["status"] == STATUS_READY # 从磁盘复核状态已持久化 assert tm.load_task("t1").state.status == STATUS_READY def test_process_with_review_goes_reviewing(isolated_env): """产生 review 桶 → 任务状态 reviewing""" # 多音字/字母混合等会进 review;用一个含字母的确保进 special lines = ["日本人:にほんじん:", "IT:アイティー:"] tm, task = _new_task(isolated_env, lines) bp = BatchProcessor(task) result = bp.process_batch() review_total = sum( result["buckets"][b] for b in ("pinyin", "split", "verb", "special") ) # skip 也可能吃掉纯字母词;只要不是全部 auto 即可能有 review if review_total > 0: assert result["status"] == STATUS_REVIEWING else: assert result["status"] == STATUS_READY # -------------------------------------------------------------------------- # 合并授权门禁 # -------------------------------------------------------------------------- def test_merge_rejected_when_review_pending(isolated_env): """review 未清零时 merge_all 必须拒绝(授权门禁的前置约束)""" tm, task = _new_task(isolated_env, ["日本人:にほんじん:"]) bp = BatchProcessor(task) bp.process_batch() # 人为往 review 桶塞一条,模拟待确认 bp.review_files["pinyin"].write_text("行|列:こう|れつ:hang|lie\n", encoding="utf-8", newline="\n") result = bp.merge_all(dry_run=True) assert "error" in result assert result["total_review"] > 0 def test_merge_final_dedup_and_status(isolated_env): """review 清零后合并到最终库:去重 + 状态变 merged + 单批清空""" lines = ["日本人:にほんじん:", "中国人:ちゅうごくじん:"] tm, task = _new_task(isolated_env, lines) bp = BatchProcessor(task) bp.process_batch() assert bp.get_status()["ready_to_merge"] is True # 首次合并 r1 = bp.merge_all(dry_run=False) assert r1["merged_auto"] == 2 assert r1["single_batch_cleared"] is True assert tm.load_task("t1").state.status == STATUS_MERGED assert r1["main_total"] == 2 def test_merge_idempotent(isolated_env): """幂等:同一批数据再次进入主库不会重复""" lines = ["日本人:にほんじん:"] tm, task = _new_task(isolated_env, lines) bp = BatchProcessor(task) bp.process_batch() bp.merge_all(dry_run=False) # 再处理同样的数据 + 合并,主库不应增长 bp.process_batch() r2 = bp.merge_all(dry_run=False) assert r2["merged_auto"] == 0 assert r2["duplicated_auto"] == 1 assert r2["main_total"] == 1 # -------------------------------------------------------------------------- # review 重跑(改代码后重新分流) # -------------------------------------------------------------------------- def test_reprocess_review_reclassifies(isolated_env): """重跑 review 桶:条目重新分类,原桶清零""" tm, task = _new_task(isolated_env, ["日本人:にほんじん:"]) bp = BatchProcessor(task) bp.process_batch() # 手动放一条可自动处理的词进 pinyin 桶,模拟"修正后应归入 auto" bp.review_files["pinyin"].write_text("中国人:ちゅうごくじん:\n", encoding="utf-8", newline="\n") result = bp.reprocess_review("pinyin") assert result["reprocessed"] == 1 # 原 pinyin 桶已清零 assert len(bp._read_clean_lines(bp.review_files["pinyin"])) == 0 def test_reprocess_unknown_bucket_raises(isolated_env): """重跑未知桶名报错""" tm, task = _new_task(isolated_env, ["日本:にほん:"]) bp = BatchProcessor(task) with pytest.raises(ValueError): bp.reprocess_review("nonexistent")