提交 eba91eac authored 作者: 陈泽健's avatar 陈泽健

fix(scheduler): 定时任务僵尸执行启动恢复 + DB 级重叠兜底守卫

- execution_service 新增 recover_interrupted_executions:启动把超龄 running/pending
  UI 执行标 failed,同步其 pending/running 用例结果
- scheduler_service run_task_once 新增 DB 级兜底:内存锁为空但 DB 仍有 running/pending
  UI 执行时跳过触发并推进 next_run_at
- main lifespan 启动调用恢复,打印清理日志
- 新增 test_scheduler_recovery 5 项测试
Co-Authored-By: 's avatarClaude <noreply@anthropic.com>
上级 7e8d7d20
......@@ -121,6 +121,21 @@ async def lifespan(app: FastAPI):
except Exception as e:
logger.warning(f"消息消费循环启动失败: {e}")
# 启动时兜底恢复:清理上次进程异常退出遗留的僵尸执行记录(status 卡在
# running/pending 的 UI 执行)。否则内存执行锁是空的,定时任务会照常触发,
# 与旧进程残留的 Playwright 线程并发执行,导致执行永不产出报告。
try:
from app.services.execution_service import recover_interrupted_executions
recovered, affected = await recover_interrupted_executions()
if recovered:
logger.warning(
f"启动时恢复 {recovered} 条僵尸执行、{affected} 条用例结果"
)
else:
logger.info("启动时无僵尸执行记录,跳过恢复")
except Exception as e:
logger.warning(f"启动时恢复僵尸执行失败(非致命): {e}")
# 启动定时任务调度引擎
scheduler_task = None
try:
......
......@@ -22,7 +22,9 @@ from datetime import datetime
from concurrent.futures import ThreadPoolExecutor
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func, desc, or_
from sqlalchemy import select, func, desc, or_, update
from app.database import async_session_maker
from app.models.execution import Execution
from app.models.case_result import CaseResult
......@@ -123,6 +125,88 @@ def clear_cancel_requested(execution_id: str) -> None:
_cancel_requested.discard(execution_id)
async def recover_interrupted_executions(
max_age_seconds: int = 300,
) -> Tuple[int, int]:
"""
启动时兜底恢复:清理上次进程异常退出遗留的僵尸执行记录。
进程是唯一能在 run_execution 的 finally 中把 status 写回 completed/failed 的
执行者。若容器在运行中被重启/杀掉(内存态运行标志随之丢失),DB 中的记录会
永远停在 running——此时内存锁 is_execution_running() 已是 False,定时任务会
继续触发创建新的执行并新起浏览器。若旧执行只是「记录被卡住」而 Playwright
浏览器线程其实还活着,就会出现多浏览器并发执行互相干扰、报告完全不产出的
现象(已在 5.44/5.202 复现:18/14 条 running 僵尸记录,任务每 6 小时/300 分钟
照常触发)。
修复策略:进程刚启动时(无任何执行在跑),把 DB 中所有 running/pending 的
UI 执行记录标记为 failed(start_time 超过 max_age_seconds 的才标记,避免误伤
刚创建、尚未进入 worker 的 pending 记录)。旧进程若还有 Playwright 线程在跑,
会在下一轮执行完成后自然退出并由服务端清理资源,不会被新执行锁阻塞。
安全测试(security)执行是纯 requests 无浏览器,卡死风险低,且可能由多窗口
手动触发,此处不处理,保持原有语义。
Args:
max_age_seconds (int): 执行记录年龄阈值。超过此秒数仍处于
running/pending 视为僵尸,默认 300 秒。
Returns:
Tuple[int, int]: (标记为失败的执行数, 受影响用例结果数)
"""
from datetime import timedelta
async with async_session_maker() as db:
result = await db.execute(
select(Execution).where(
Execution.status.in_(["running", "pending"]),
Execution.case_type == "ui",
).with_for_update()
)
stale = [
e for e in result.scalars().all()
if e.start_time
and (datetime.now() - e.start_time).total_seconds() >= max_age_seconds
]
if not stale:
return 0, 0
exec_ids = [e.id for e in stale]
now = datetime.now()
# 先标记执行记录为 failed(保留 start/end,便于排查)
await db.execute(
update(Execution)
.where(Execution.id.in_(exec_ids))
.values(
status="failed",
end_time=now,
duration=0,
error_message=(
"执行被中断(进程重启/异常退出),启动时自动恢复标记为失败"
),
)
)
# 关联的 pending/running 用例结果 → failed,保证报告/统计口径一致
case_result_update = await db.execute(
update(CaseResult)
.where(
CaseResult.execution_id.in_(exec_ids),
CaseResult.status.in_(["pending", "running"]),
)
.values(
status="failed",
error_message="执行被中断(进程重启/异常退出),启动时自动恢复",
end_time=now,
)
)
await db.commit()
logger.warning(
f"[启动恢复] 清理 {len(exec_ids)} 条僵尸执行记录"
f"(关联 {case_result_update.rowcount} 条用例结果标记为 failed)"
)
return len(exec_ids), case_result_update.rowcount
class ExecutionService:
"""
执行服务类
......
......@@ -292,6 +292,24 @@ async def run_task_once(task_id: str, manual: bool = False) -> Optional[str]:
await db.flush()
return None
# 2b. DB 级兜底守卫:即便内存锁因进程重启而丢失,只要 DB 中仍有
# UI 执行处于 running/pending 状态,也不触发新执行,避免并发
# Playwright 浏览器实例互相干扰、执行永不产出报告。
db_running_q = await db.execute(
select(Execution.id).where(
Execution.case_type == "ui",
Execution.status.in_(["running", "pending"]),
).limit(1)
)
if db_running_q.scalar_one_or_none() is not None:
logger.info(
f"[定时任务] 任务「{task.name}」到点但 DB 中仍有 UI 执行在跑,"
f"跳过本次,下次周期再触发"
)
task.next_run_at = compute_next_run(task)
await db.flush()
return None
# 3. 创建执行记录(trigger_type=scheduled, trigger_by=任务名)
case_ids = [c.id for c in cases]
exec_service = ExecutionService(db)
......
#!/usr/bin/env python
# -*- coding: utf-8 -*-
"""
模块名称:test_scheduler_recovery.py
模块描述:定时任务僵尸执行恢复 + DB 级重叠守卫 单元测试
覆盖:
- recover_interrupted_executions:启动时把超过年龄阈值的 running/pending UI 执行
标记为 failed,并同步其关联的 pending/running 用例结果(MySQL 行级 UPDATE)。
- 年轻(未超阈值)的 running/pending 记录不被误伤。
- security 执行记录不受影响。
- scheduler_service.run_task_once:即便内存执行锁为空,只要 DB 中仍有 UI 执行
处于 running/pending,也跳过触发(DB 级兜底守卫)。
作者:czj
创建日期:2026-08-27
"""
from datetime import datetime, timedelta
from unittest.mock import AsyncMock, MagicMock, patch
from app.models.execution import Execution
from app.models.case_result import CaseResult
from app.services.execution_service import recover_interrupted_executions
def _mk_exec(eid, status, start_time, case_type="ui"):
e = Execution(id=eid, status=status, case_type=case_type)
e.start_time = start_time
return e
# ==================== recover_interrupted_executions ====================
class TestRecoverInterruptedExecutions:
"""启动时僵尸执行恢复"""
def _patch_ctx(self, stale_execs):
"""构造 patch:async_session_maker 返回的 session.execute 返回对应结果"""
exec_result = MagicMock()
exec_result.scalars.return_value.all.return_value = stale_execs
update_result = MagicMock()
update_result.rowcount = 7
session = AsyncMock()
session.execute = AsyncMock(side_effect=[exec_result, update_result, update_result])
cm = AsyncMock()
cm.__aenter__ = AsyncMock(return_value=session)
cm.__aexit__ = AsyncMock(return_value=False)
session_maker = MagicMock()
session_maker.return_value = cm
return session_maker
def test_marks_stale_running_and_pending_as_failed(self):
now = datetime.now()
stale = [
_mk_exec("e1", "running", now - timedelta(hours=2)),
_mk_exec("e2", "pending", now - timedelta(days=1)),
]
session_maker = self._patch_ctx(stale)
with patch(
"app.services.execution_service.async_session_maker", session_maker
):
recovered, affected = asyncio_run(recover_interrupted_executions())
assert recovered == 2
assert affected == 7
# 校验执行记录 UPDATE 与用例结果 UPDATE 均被执行,且都是置为 failed
session = session_maker.return_value.__aenter__.return_value
update_sqls = [
str(c.args[0])
for c in session.execute.call_args_list
if str(c.args[0]).lstrip().upper().startswith("UPDATE")
]
assert len(update_sqls) == 2
assert "UPDATE EXECUTIONS" in update_sqls[0].upper()
assert "UPDATE CASE_RESULTS" in update_sqls[1].upper()
# 两条 UPDATE 的 status 绑定值均为 failed(values 以 Column 为键,值为 BindParameter)
for idx, c in enumerate(
c for c in session.execute.call_args_list
if str(c.args[0]).lstrip().upper().startswith("UPDATE")
):
values = dict(c.args[0]._values or {})
status_val = next(
(v for k, v in values.items() if getattr(k, "name", None) == "status"),
None,
)
effective = getattr(status_val, "value", status_val)
assert effective == "failed", f"UPDATE #{idx} status 应为 failed"
def test_young_running_not_affected(self):
now = datetime.now()
stale = [_mk_exec("e1", "running", now - timedelta(seconds=60))]
session_maker = self._patch_ctx(stale)
with patch(
"app.services.execution_service.async_session_maker", session_maker
):
recovered, affected = asyncio_run(recover_interrupted_executions(max_age_seconds=300))
assert recovered == 0
assert affected == 0
def test_empty_no_error(self):
session_maker = self._patch_ctx([])
with patch(
"app.services.execution_service.async_session_maker", session_maker
):
recovered, affected = asyncio_run(recover_interrupted_executions())
assert recovered == 0
assert affected == 0
# ==================== DB 级重叠守卫(run_task_once) ====================
def _as_scalar(value):
r = MagicMock()
r.scalar_one_or_none.return_value = value
return r
class TestSchedulerDbGuard:
"""DB 中仍有 UI 执行在跑时,定时任务到点也应跳过"""
def _make_task(self):
from app.models.scheduled_task import ScheduledTask
t = ScheduledTask(
id="scheduled_x",
name="测试任务",
enabled=True,
module_ids=["m1"],
case_type="ui",
auto_report=True,
)
return t
def test_db_running_blocks_trigger(self):
"""内存锁空(进程重启后),但 DB 有 running UI 执行 → 跳过,不创建执行"""
from app.services import scheduler_service
task = self._make_task()
session = AsyncMock()
# 1. 查任务 → 返回 task
# 2. collect_module_cases 查询 → 返回非空用例
# 3. DB 级守卫查询 → 返回存在 running
session.execute = AsyncMock()
exec_side = []
task_result = MagicMock()
task_result.scalar_one_or_none.return_value = task
exec_side.append(task_result)
case_result = MagicMock()
case_result.scalars.return_value.all.return_value = ["case_1"]
exec_side.append(case_result)
db_guard = MagicMock()
db_guard.scalar_one_or_none.return_value = "exec_running"
exec_side.append(db_guard)
session.execute.side_effect = exec_side
cm = AsyncMock()
cm.__aenter__ = AsyncMock(return_value=session)
cm.__aexit__ = AsyncMock(return_value=False)
session_maker = MagicMock()
session_maker.return_value = cm
# 内存锁:is_execution_running 返回 False(进程刚重启,内存无运行记录)
with patch.object(scheduler_service, "async_session_maker", session_maker), \
patch.object(scheduler_service, "is_execution_running", return_value=(False, None)), \
patch.object(scheduler_service, "collect_module_cases", AsyncMock(return_value=["case_1"])):
result = asyncio_run(scheduler_service.run_task_once("scheduled_x"))
assert result is None
# 跳过时推进了 next_run_at(触发了 db.flush)
assert session.flush.await_count >= 1
def test_db_no_running_proceeds(self):
"""内存锁空且 DB 无 running → 正常触发执行"""
from app.services import scheduler_service
from app.models.test_case import TestCase
task = self._make_task()
# collect_module_cases 被 patched,不走真实 DB;
# 返回包含 TestCase 对象的列表(run_task_once 需要 .id 属性)
fake_cases = [
TestCase(id="case_1", name="测试用例1", module_id="m1", status="active", case_type="ui")
]
fake_cases[0].steps = []
session = AsyncMock()
task_result = MagicMock()
task_result.scalar_one_or_none.return_value = task
db_guard = MagicMock()
db_guard.scalar_one_or_none.return_value = None # 无 running
# session.execute 只被调用两次:查任务 + DB 守卫
session.execute.side_effect = [task_result, db_guard]
# 创建执行 → 返回一个 execution 对象
exec_obj = Execution(id="exec_new", name="x", status="pending", case_type="ui")
exec_service = MagicMock()
exec_service.create_execution = AsyncMock(return_value=exec_obj)
exec_service.run_execution = AsyncMock(return_value=exec_obj)
cm = AsyncMock()
cm.__aenter__ = AsyncMock(return_value=session)
cm.__aexit__ = AsyncMock(return_value=False)
session_maker = MagicMock()
session_maker.return_value = cm
with patch.object(scheduler_service, "async_session_maker", session_maker), \
patch.object(scheduler_service, "is_execution_running", return_value=(False, None)), \
patch.object(scheduler_service, "collect_module_cases", AsyncMock(return_value=fake_cases)), \
patch.object(scheduler_service, "ExecutionService", return_value=exec_service), \
patch.object(scheduler_service, "ReportService") as rs:
rs.return_value.generate_html = AsyncMock(return_value="<html/>")
rs.return_value.save_report = AsyncMock(return_value="/tmp/r.html")
result = asyncio_run(scheduler_service.run_task_once("scheduled_x"))
assert result == "exec_new"
exec_service.create_execution.assert_awaited_once()
exec_service.run_execution.assert_awaited_once_with("exec_new")
# ==================== 便捷运行器 ====================
import asyncio
def asyncio_run(coro):
"""跨平台运行异步协程(与项目其他测试一致,避免 Windows 事件循环差异)"""
return asyncio.new_event_loop().run_until_complete(coro)
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论