ForcePilot/backend/test/integration/api/test_scheduler_router.py
Kris e9a1ea1050 test(scheduler): add complete integration tests for scheduler router
覆盖定时任务调度路由的所有核心场景:三级权限校验、CRUD全流程、状态机操作、回收站管理、各类查询端点、参数校验以及运维接口,对齐现有测试规范实现
2026-06-23 13:42:15 +08:00

535 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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