"""任务级持久化测试(T16,OV7)。 OV7 裁定:崩溃恢复只到会话级,任务层丢数据 → 任务级持久化。 实现:api-design §5 TaskQueue 抽象 + PersistentTaskQueue(SQLite 落盘), 重启后可 recover 未完成任务(running → failed,pending 保留)。 """ from __future__ import annotations from genesis.orchestrator.task_queue import ( PersistentTaskQueue, TaskSpec, TaskStatus, ) def _spec(**overrides) -> TaskSpec: base = dict( task_id="t-1", session_id="s-1", step="generate", chapter_id="ch-3", idempotency_key="s-1|generate|ch-3", payload={"chapter_id": "ch-3"}, ) base.update(overrides) return TaskSpec(**base) def test_enqueue_and_get(tmp_path): q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") handle = q.enqueue(_spec()) got = q.get("t-1") assert got is not None assert got.status == TaskStatus.PENDING assert got.payload == {"chapter_id": "ch-3"} q.close() def test_poll_returns_session_tasks(tmp_path): q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") q.enqueue(_spec(task_id="t-1")) q.enqueue(_spec(task_id="t-2", chapter_id="ch-4", idempotency_key="s-1|generate|ch-4")) tasks = q.poll("s-1") assert {t.task_id for t in tasks} == {"t-1", "t-2"} q.close() def test_update_status_with_result(tmp_path): q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") q.enqueue(_spec()) q.update_status("t-1", TaskStatus.COMPLETED, result={"html": "

"}) got = q.get("t-1") assert got.status == TaskStatus.COMPLETED assert got.result == {"html": "

"} q.close() def test_cancel(tmp_path): q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") q.enqueue(_spec()) assert q.cancel("t-1") is True assert q.get("t-1").status == TaskStatus.CANCELLED assert q.cancel("t-1") is False # 已终态不可再取消 q.close() def test_idempotent_enqueue_returns_cached_result(tmp_path): """§5.3 幂等去重:同 idempotency_key 已完成 → 返回缓存结果,不重复入队。""" q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") q.enqueue(_spec()) q.update_status("t-1", TaskStatus.COMPLETED, result={"html": "cached"}) second = q.enqueue(_spec(task_id="t-999")) assert second.task_id == "t-1" # 返回原 handle assert second.result == {"html": "cached"} assert q.get("t-999") is None q.close() def test_persistence_survives_reopen(tmp_path): """OV7 核心:任务写入 SQLite,重启(重建实例)后数据不丢。""" db = tmp_path / "tasks.db" q1 = PersistentTaskQueue(db_path=db) q1.enqueue(_spec()) q1.update_status("t-1", TaskStatus.RUNNING) q1.close() q2 = PersistentTaskQueue(db_path=db) got = q2.get("t-1") assert got is not None assert got.status == TaskStatus.RUNNING q2.close() def test_recover_marks_running_as_failed_keeps_pending(tmp_path): """崩溃恢复:running → failed(中断标记),pending 保留待执行,completed 不动。""" db = tmp_path / "tasks.db" q1 = PersistentTaskQueue(db_path=db) q1.enqueue(_spec(task_id="t-running")) q1.update_status("t-running", TaskStatus.RUNNING) q1.enqueue(_spec(task_id="t-pending", chapter_id="ch-5", idempotency_key="s-1|generate|ch-5")) q1.enqueue(_spec(task_id="t-done", chapter_id="ch-6", idempotency_key="s-1|generate|ch-6")) q1.update_status("t-done", TaskStatus.COMPLETED, result={"html": "ok"}) q1.close() q2 = PersistentTaskQueue(db_path=db) recovered = q2.recover() assert q2.get("t-running").status == TaskStatus.FAILED assert q2.get("t-pending").status == TaskStatus.PENDING assert q2.get("t-done").status == TaskStatus.COMPLETED assert {t.task_id for t in recovered} == {"t-running", "t-pending"} # 未完成待处理 q2.close() def test_poll_only_incomplete_after_recover(tmp_path): db = tmp_path / "tasks.db" q1 = PersistentTaskQueue(db_path=db) q1.enqueue(_spec(task_id="t-running")) q1.update_status("t-running", TaskStatus.RUNNING) q1.close() q2 = PersistentTaskQueue(db_path=db) q2.recover() pending = q2.poll("s-1") assert {t.task_id for t in pending} == {"t-running"} # 已标记 failed,仍可重试 q2.close() # ---------- 防御分支(覆盖率 100% 基线) ---------- def test_update_status_unknown_task_raises(tmp_path): import pytest q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") with pytest.raises(KeyError, match="t-404"): q.update_status("t-404", TaskStatus.RUNNING) q.close() def test_update_status_terminal_rejected(tmp_path): import pytest q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") q.enqueue(_spec()) q.update_status("t-1", TaskStatus.COMPLETED, result={"html": "done"}) with pytest.raises(ValueError, match="终态"): q.update_status("t-1", TaskStatus.RUNNING) q.close() def test_idem_lookup_with_null_chapter_id(tmp_path): """幂等查找:chapter_id 为 None(非章级任务)时仍命中已有任务。""" q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") spec = _spec(chapter_id=None, task_id="t-impact", step="impact", idempotency_key="s-1|impact|") q.enqueue(spec) second = q.enqueue(_spec(chapter_id=None, task_id="t-impact-2", step="impact", idempotency_key="s-1|impact|")) assert second.task_id == "t-impact" # 命中已有,未重复入队 q.close() def test_row_to_handle_unknown_task_raises(tmp_path): """白盒:_row_to_handle 无行时抛 KeyError(内部防御分支)。""" import pytest q = PersistentTaskQueue(db_path=tmp_path / "tasks.db") with pytest.raises(KeyError, match="t-404"): q._row_to_handle("t-404") q.close()