test(scheduler): add complete integration tests for scheduler router
覆盖定时任务调度路由的所有核心场景:三级权限校验、CRUD全流程、状态机操作、回收站管理、各类查询端点、参数校验以及运维接口,对齐现有测试规范实现
This commit is contained in:
parent
2708ff1088
commit
e9a1ea1050
534
backend/test/integration/api/test_scheduler_router.py
Normal file
534
backend/test/integration/api/test_scheduler_router.py
Normal file
@ -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
|
||||
Loading…
Reference in New Issue
Block a user