From e9a1ea1050b04796d72d38938ec9e5955989de78 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Tue, 23 Jun 2026 13:42:15 +0800 Subject: [PATCH] test(scheduler): add complete integration tests for scheduler router MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 覆盖定时任务调度路由的所有核心场景:三级权限校验、CRUD全流程、状态机操作、回收站管理、各类查询端点、参数校验以及运维接口,对齐现有测试规范实现 --- .../integration/api/test_scheduler_router.py | 534 ++++++++++++++++++ 1 file changed, 534 insertions(+) create mode 100644 backend/test/integration/api/test_scheduler_router.py diff --git a/backend/test/integration/api/test_scheduler_router.py b/backend/test/integration/api/test_scheduler_router.py new file mode 100644 index 00000000..75b449fd --- /dev/null +++ b/backend/test/integration/api/test_scheduler_router.py @@ -0,0 +1,534 @@ +"""Integration tests for scheduler_router endpoints. + +覆盖定时任务调度限界上下文 Router(/api/scheduler/*):任务 CRUD / 状态机 / +手动触发 / 执行日志查询 / 可观测性 / 健康检查 / 运维恢复(PRD §FR-ST-01 ~ §FR-ST-09)。 + +对齐 external_systems/test_system_router.py 写法: +- 三级认证(401 未认证 / 403 普通用户 / 200 管理员) +- CRUD 全流程(创建 → 查询 → 更新 → 删除) +- 状态机(pause / resume / restore / hard delete) +- 查询端点(upcoming / count-by-status / run-logs / handlers/summary / stats/daily / health) +- 参数校验(422) +""" + +from __future__ import annotations + +import uuid +from collections.abc import AsyncGenerator +from datetime import UTC, datetime, timedelta + +import httpx +import pytest +import pytest_asyncio + +pytestmark = [pytest.mark.asyncio, pytest.mark.integration] + +BASE_URL = "/api/scheduler" +TASKS_URL = f"{BASE_URL}/tasks" + + +# ============================================================================= +# === Helpers & fixtures === +# ============================================================================= + + +def _make_cron_payload() -> dict: + """构造合法的 cron 周期任务请求体。handler_name 仅作引用标识,api 进程不校验存在性。""" + suffix = uuid.uuid4().hex[:8] + return { + "handler_name": f"pytest_handler_{suffix}", + "owner_scope": "system", + "owner_id": f"pytest:{suffix}", + "schedule_kind": "cron", + "cron_expression": "0 * * * *", + "tz": "Asia/Shanghai", + "payload": {"key": "value"}, + "enabled": True, + } + + +def _make_at_payload() -> dict: + """构造合法的 at 一次性任务请求体,run_at 取未来 1 天。""" + suffix = uuid.uuid4().hex[:8] + future = (datetime.now(UTC) + timedelta(days=1)).isoformat() + return { + "handler_name": f"pytest_handler_{suffix}", + "owner_scope": "system", + "owner_id": f"pytest:{suffix}", + "schedule_kind": "at", + "run_at": future, + "tz": "Asia/Shanghai", + "payload": {}, + "enabled": True, + } + + +@pytest_asyncio.fixture(scope="function") +async def scheduler_task( + test_client: httpx.AsyncClient, + admin_headers: dict[str, str], +) -> AsyncGenerator[dict, None]: + """创建一个 cron 周期任务供测试复用,用例结束后硬删除清理。""" + response = await test_client.post(TASKS_URL, json=_make_cron_payload(), headers=admin_headers) + assert response.status_code == 200, f"Failed to create scheduler task: {response.text}" + task = response.json()["data"] + try: + yield task + finally: + # 先软删除再硬删除,覆盖 active / deleted 两种状态 + for url in (f"{TASKS_URL}/{task['task_id']}", f"{TASKS_URL}/{task['task_id']}/hard"): + try: + await test_client.delete(url, headers=admin_headers) + except Exception: + pass + + +# ============================================================================= +# === Auth three-tier for list endpoint === +# ============================================================================= + + +async def test_list_tasks_requires_auth(test_client): + # Act + response = await test_client.get(TASKS_URL) + # Assert + assert response.status_code == 401 + + +async def test_list_tasks_allows_standard_user(test_client, standard_user): + # Act + response = await test_client.get(TASKS_URL, headers=standard_user["headers"]) + # Assert + assert response.status_code == 200, response.text + + +async def test_admin_can_list_tasks(test_client, admin_headers): + # Act + response = await test_client.get(TASKS_URL, headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + assert isinstance(payload["data"]["items"], list) + assert payload["data"]["total"] == len(payload["data"]["items"]) + + +# ============================================================================= +# === Auth for create endpoint === +# ============================================================================= + + +async def test_create_task_requires_auth(test_client): + # Act + response = await test_client.post(TASKS_URL, json=_make_cron_payload()) + # Assert + assert response.status_code == 401 + + +async def test_create_task_requires_admin(test_client, standard_user): + # Act + response = await test_client.post(TASKS_URL, json=_make_cron_payload(), headers=standard_user["headers"]) + # Assert + assert response.status_code == 403 + + +async def test_admin_can_create_cron_task(test_client, admin_headers, scheduler_task): + # Assert: fixture 已完成创建,校验关键字段 + assert scheduler_task["task_id"] + assert scheduler_task["schedule_kind"] == "cron" + assert scheduler_task["status"] == "active" + assert scheduler_task["next_run_at"] is not None + assert scheduler_task["stagger_seconds"] is not None + + +async def test_admin_can_create_at_task(test_client, admin_headers): + # Arrange + payload = _make_at_payload() + # Act + response = await test_client.post(TASKS_URL, json=payload, headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + data = response.json()["data"] + assert data["schedule_kind"] == "at" + assert data["run_at"] is not None + assert data["next_run_at"] is not None + # Cleanup + try: + await test_client.delete(f"{TASKS_URL}/{data['task_id']}/hard", headers=admin_headers) + except Exception: + pass + + +# ============================================================================= +# === CRUD flow === +# ============================================================================= + + +async def test_admin_can_get_task_detail(test_client, admin_headers, scheduler_task): + # Act + response = await test_client.get(f"{TASKS_URL}/{scheduler_task['task_id']}", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + data = response.json()["data"] + assert data["task_id"] == scheduler_task["task_id"] + assert data["handler_name"] == scheduler_task["handler_name"] + + +async def test_get_unknown_task_returns_404(test_client, admin_headers): + # Act + response = await test_client.get(f"{TASKS_URL}/nonexistent_task_id", headers=admin_headers) + # Assert + assert response.status_code == 404, response.text + + +async def test_admin_can_update_task(test_client, admin_headers, scheduler_task): + # Arrange: 部分更新 enabled 与 payload + # Act + response = await test_client.put( + f"{TASKS_URL}/{scheduler_task['task_id']}", + json={"enabled": False, "payload": {"updated": True}}, + headers=admin_headers, + ) + # Assert + assert response.status_code == 200, response.text + data = response.json()["data"] + assert data["enabled"] is False + assert data["payload"] == {"updated": True} + assert data["task_id"] == scheduler_task["task_id"] + + +async def test_admin_can_soft_delete_task(test_client, admin_headers): + # Arrange: 内联创建以隔离删除验证 + create_response = await test_client.post(TASKS_URL, json=_make_cron_payload(), headers=admin_headers) + assert create_response.status_code == 200, create_response.text + task_id = create_response.json()["data"]["task_id"] + + try: + # Act + delete_response = await test_client.delete(f"{TASKS_URL}/{task_id}", headers=admin_headers) + # Assert + assert delete_response.status_code == 200, delete_response.text + assert delete_response.json()["data"]["status"] == "deleted" + + # 软删除后 get 应返回 404 + get_response = await test_client.get(f"{TASKS_URL}/{task_id}", headers=admin_headers) + assert get_response.status_code == 404 + finally: + try: + await test_client.delete(f"{TASKS_URL}/{task_id}/hard", headers=admin_headers) + except Exception: + pass + + +# ============================================================================= +# === State machine: pause / resume === +# ============================================================================= + + +async def test_admin_can_pause_and_resume_task(test_client, admin_headers, scheduler_task): + # Act: pause + pause_response = await test_client.post( + f"{TASKS_URL}/{scheduler_task['task_id']}/pause", headers=admin_headers + ) + # Assert + assert pause_response.status_code == 200, pause_response.text + assert pause_response.json()["data"]["status"] == "paused" + + # Act: resume + resume_response = await test_client.post( + f"{TASKS_URL}/{scheduler_task['task_id']}/resume", headers=admin_headers + ) + # Assert + assert resume_response.status_code == 200, resume_response.text + assert resume_response.json()["data"]["status"] == "active" + + +async def test_pause_requires_admin(test_client, standard_user, scheduler_task): + # Act + response = await test_client.post( + f"{TASKS_URL}/{scheduler_task['task_id']}/pause", headers=standard_user["headers"] + ) + # Assert + assert response.status_code == 403 + + +async def test_resume_active_task_returns_409(test_client, admin_headers, scheduler_task): + # Act: active 状态下 resume 不允许,抛 TaskStatusTransitionError (409) + response = await test_client.post( + f"{TASKS_URL}/{scheduler_task['task_id']}/resume", headers=admin_headers + ) + # Assert + assert response.status_code == 409, response.text + + +# ============================================================================= +# === Recycle bin: list deleted / restore / hard delete === +# ============================================================================= + + +async def test_admin_can_list_deleted_tasks(test_client, admin_headers): + # Act + response = await test_client.get(f"{TASKS_URL}/deleted", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"]["items"], list) + + +async def test_list_deleted_tasks_requires_admin(test_client, standard_user): + # Act + response = await test_client.get(f"{TASKS_URL}/deleted", headers=standard_user["headers"]) + # Assert + assert response.status_code == 403 + + +async def test_admin_can_restore_deleted_task(test_client, admin_headers): + # Arrange: 创建 → 软删除 → restore + create_response = await test_client.post(TASKS_URL, json=_make_cron_payload(), headers=admin_headers) + assert create_response.status_code == 200, create_response.text + task_id = create_response.json()["data"]["task_id"] + + try: + delete_response = await test_client.delete(f"{TASKS_URL}/{task_id}", headers=admin_headers) + assert delete_response.status_code == 200 + + # Act: restore + restore_response = await test_client.post(f"{TASKS_URL}/{task_id}/restore", headers=admin_headers) + # Assert: 恢复后 status 应为 paused(避免立即被 tick 执行) + assert restore_response.status_code == 200, restore_response.text + assert restore_response.json()["data"]["status"] == "paused" + assert restore_response.json()["data"]["is_deleted"] == 0 + + # 恢复后 get 应返回 200 + get_response = await test_client.get(f"{TASKS_URL}/{task_id}", headers=admin_headers) + assert get_response.status_code == 200 + finally: + try: + await test_client.delete(f"{TASKS_URL}/{task_id}", headers=admin_headers) + await test_client.delete(f"{TASKS_URL}/{task_id}/hard", headers=admin_headers) + except Exception: + pass + + +async def test_admin_can_hard_delete_task(test_client, admin_headers): + # Arrange: 创建 → 软删除 → 硬删除 + create_response = await test_client.post(TASKS_URL, json=_make_cron_payload(), headers=admin_headers) + assert create_response.status_code == 200, create_response.text + task_id = create_response.json()["data"]["task_id"] + + delete_response = await test_client.delete(f"{TASKS_URL}/{task_id}", headers=admin_headers) + assert delete_response.status_code == 200 + + # Act: hard delete + response = await test_client.delete(f"{TASKS_URL}/{task_id}/hard", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + assert response.json()["data"]["status"] == "hard_deleted" + + # 硬删除后 get 应返回 404 + get_response = await test_client.get(f"{TASKS_URL}/{task_id}", headers=admin_headers) + assert get_response.status_code == 404 + + +async def test_hard_delete_requires_admin(test_client, standard_user, scheduler_task): + # Act + response = await test_client.delete( + f"{TASKS_URL}/{scheduler_task['task_id']}/hard", headers=standard_user["headers"] + ) + # Assert + assert response.status_code == 403 + + +# ============================================================================= +# === Query endpoints === +# ============================================================================= + + +async def test_admin_can_list_upcoming(test_client, admin_headers): + # Act + response = await test_client.get(f"{TASKS_URL}/upcoming", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + + +async def test_admin_can_count_by_status(test_client, admin_headers): + # Act + response = await test_client.get(f"{TASKS_URL}/count-by-status", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + + +async def test_admin_can_list_run_logs(test_client, admin_headers, scheduler_task): + # Act + response = await test_client.get( + f"{TASKS_URL}/{scheduler_task['task_id']}/run-logs", headers=admin_headers + ) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"]["items"], list) + + +async def test_admin_can_list_all_run_logs(test_client, admin_headers): + # Act: 跨任务执行日志查询,task_id 与 status 至少传一个(用例层校验) + response = await test_client.get(f"{BASE_URL}/run-logs?status=success", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"]["items"], list) + + +async def test_list_all_run_logs_rejects_no_filter(test_client, admin_headers): + # Act: task_id 与 status 均为空时应被用例层拒绝(400) + response = await test_client.get(f"{BASE_URL}/run-logs", headers=admin_headers) + # Assert + assert response.status_code == 400, response.text + + +async def test_admin_can_get_run_log_unknown_returns_404(test_client, admin_headers): + # Act + response = await test_client.get(f"{BASE_URL}/run-logs/nonexistent_run_id", headers=admin_headers) + # Assert + assert response.status_code == 404, response.text + + +async def test_admin_can_list_handler_summary(test_client, admin_headers): + # Act + response = await test_client.get(f"{BASE_URL}/handlers/summary", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], list) + + +async def test_admin_can_list_daily_stats(test_client, admin_headers): + # Arrange: 取近 7 天范围 + end = datetime.now(UTC).strftime("%Y-%m-%d") + start = (datetime.now(UTC) - timedelta(days=7)).strftime("%Y-%m-%d") + # Act + response = await test_client.get( + f"{BASE_URL}/stats/daily?start_date={start}&end_date={end}", headers=admin_headers + ) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + + +async def test_daily_stats_requires_date_params(test_client, admin_headers): + # Act: start_date / end_date 为必填,缺失应返回 422 + response = await test_client.get(f"{BASE_URL}/stats/daily", headers=admin_headers) + # Assert + assert response.status_code == 422, response.text + + +async def test_admin_can_get_health(test_client, admin_headers): + # Act + response = await test_client.get(f"{BASE_URL}/health", headers=admin_headers) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + + +# ============================================================================= +# === Ops: reclaim stale runs === +# ============================================================================= + + +async def test_admin_can_reclaim_stale_runs(test_client, admin_headers): + # Act + response = await test_client.post( + f"{BASE_URL}/reclaim-stale-runs", json={"timeout_seconds": 3600}, headers=admin_headers + ) + # Assert + assert response.status_code == 200, response.text + payload = response.json() + assert payload["success"] is True + assert isinstance(payload["data"], dict) + + +async def test_reclaim_stale_runs_requires_admin(test_client, standard_user): + # Act + response = await test_client.post( + f"{BASE_URL}/reclaim-stale-runs", headers=standard_user["headers"] + ) + # Assert + assert response.status_code == 403 + + +# ============================================================================= +# === Parameter validation === +# ============================================================================= + + +async def test_create_task_rejects_invalid_handler_name(test_client, admin_headers): + # Arrange: handler_name 以数字开头违反 pattern + payload = _make_cron_payload() + payload["handler_name"] = "1invalid" + # Act + response = await test_client.post(TASKS_URL, json=payload, headers=admin_headers) + # Assert + assert response.status_code == 422, response.text + + +async def test_create_task_rejects_invalid_schedule_kind(test_client, admin_headers): + # Arrange + payload = _make_cron_payload() + payload["schedule_kind"] = "invalid" + # Act + response = await test_client.post(TASKS_URL, json=payload, headers=admin_headers) + # Assert + assert response.status_code == 422, response.text + + +async def test_create_cron_task_without_expression_returns_400(test_client, admin_headers): + # Arrange: schedule_kind=cron 但未传 cron_expression,业务校验失败 + payload = _make_cron_payload() + payload["cron_expression"] = None + # Act + response = await test_client.post(TASKS_URL, json=payload, headers=admin_headers) + # Assert + assert response.status_code == 400, response.text + + +async def test_create_task_rejects_invalid_cron_expression(test_client, admin_headers): + # Arrange: cron 表达式语法非法 + payload = _make_cron_payload() + payload["cron_expression"] = "not a cron" + # Act + response = await test_client.post(TASKS_URL, json=payload, headers=admin_headers) + # Assert + assert response.status_code == 400, response.text + + +async def test_list_tasks_rejects_invalid_pagination(test_client, admin_headers): + # Act: page=0 违反 ge=1 + response = await test_client.get(f"{TASKS_URL}?page=0", headers=admin_headers) + # Assert + assert response.status_code == 422, response.text + + # Act: page_size=101 违反 le=100 + response = await test_client.get(f"{TASKS_URL}?page_size=101", headers=admin_headers) + # Assert + assert response.status_code == 422, response.text + + +async def test_list_tasks_rejects_invalid_status_filter(test_client, admin_headers): + # Act: status 不在枚举内 + response = await test_client.get(f"{TASKS_URL}?status=invalid", headers=admin_headers) + # Assert + assert response.status_code == 422, response.text