diff --git a/backend/app/api/v1/daily_reports.py b/backend/app/api/v1/daily_reports.py index 7a2d9754..388a6ec0 100644 --- a/backend/app/api/v1/daily_reports.py +++ b/backend/app/api/v1/daily_reports.py @@ -227,7 +227,6 @@ async def trigger_generate_version( 改为异步:先创建/标记 GENERATING 记录立即返回 202, 后台 task 完成后前端通过轮询 /today 拿最终结果。 """ - import asyncio from app.services.daily_report import _day_window, _local_today, _local_window_to_utc_naive @@ -284,14 +283,17 @@ async def _bg_generate(): ) except Exception as e: logger.error("Background daily report generation failed: %s", e) - # 标记失败 + # 标记失败;mark_error 再失败必须留痕(曾因 except: pass + # 导致报告永远卡在 GENERATING,issue #72),日志兜底。 try: bg_repo = DailyReportRepository(bg_db) await bg_repo.mark_error(report_id, f"生成失败: {str(e)[:200]}") except Exception: - pass + logger.exception("Failed to mark daily report id=%s as ERROR", report_id) - asyncio.create_task(_bg_generate()) + from app.core.task_registry import track_background_task + + track_background_task(_bg_generate(), name=f"daily-report-bg-{report_id}") # 返回 GENERATING 状态(HTTP 202 Accepted) from fastapi.responses import JSONResponse diff --git a/backend/app/core/task_registry.py b/backend/app/core/task_registry.py new file mode 100644 index 00000000..1ef095e4 --- /dev/null +++ b/backend/app/core/task_registry.py @@ -0,0 +1,53 @@ +"""受管后台任务注册表。 + +fire-and-forget 的 `create_task` 若不保存引用、不收集异常,会静默丢任务 +(issue #72),且可能在优雅停机后仍与 `engine.dispose()` 竞争打开中的 +session。本模块集中提供: + +- `track_background_task(coro, *, name)`:创建任务并注册,done callback + 自动注销并在异常时记 error 日志(CancelledError 除外); +- `drain_tracked_tasks(timeout)`:停机时统一 cancel + 等待,超时的残留 + 任务记 warning 后放行(不强杀,避免打断不可取消的清理逻辑)。 + +兴趣向量重建(`interest_vector_service._rebuild_tasks`)是仓内同模式的 +既有范例,本模块将其通用化。 +""" + +from __future__ import annotations + +import asyncio +import logging +from collections.abc import Coroutine + +logger = logging.getLogger(__name__) + +_tasks: set[asyncio.Task] = set() + + +def _detach(task: asyncio.Task) -> None: + _tasks.discard(task) + if task.cancelled(): + return + exc = task.exception() + if exc is not None: + logger.error("Background task %s failed: %s", task.get_name(), exc, exc_info=exc) + + +def track_background_task(coro: Coroutine, *, name: str) -> asyncio.Task: + """创建并注册受管后台任务;异常由 done callback 记日志,不再静默。""" + task = asyncio.create_task(coro, name=name) + _tasks.add(task) + task.add_done_callback(_detach) + return task + + +async def drain_tracked_tasks(timeout: float = 30.0) -> None: + """优雅停机:取消全部受管任务并在 timeout 内等待结束。""" + if not _tasks: + return + for task in list(_tasks): + task.cancel() + done, pending = await asyncio.wait(set(_tasks), timeout=timeout) + for task in pending: + logger.warning("Background task %s outlived shutdown drain (timeout=%.0fs)", task.get_name(), timeout) + _tasks.clear() diff --git a/backend/app/main.py b/backend/app/main.py index 9f2e5513..c873523b 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -424,6 +424,12 @@ async def _seed_model_catalog() -> None: await _jieba_prewarm_task shutdown_scheduler() + # 回收受管后台任务(调度器启动 rescan/恢复、日报后台生成等), + # 避免 loop 关闭腰斩或与 engine.dispose 竞争 session(issue #72)。 + from app.core.task_registry import drain_tracked_tasks + + await drain_tracked_tasks(timeout=30.0) + # Drain interest-vector background rebuild tasks try: from app.services.interest_vector_service import drain_rebuild_tasks diff --git a/backend/app/scheduler.py b/backend/app/scheduler.py index 9e77b2c0..8d5c1d3f 100644 --- a/backend/app/scheduler.py +++ b/backend/app/scheduler.py @@ -1196,14 +1196,15 @@ async def _cleanup_old_metrics_snapshots() -> None: # Immediately register all enabled sources so they start syncing # right away instead of waiting for the first 10-minute rescan. - import asyncio as _asyncio + # 启动任务经 track_background_task 受管:异常由 done callback 记日志, + # 优雅停机时由 drain_tracked_tasks 统一回收(issue #72)。 + from app.core.task_registry import track_background_task try: - loop = _asyncio.get_event_loop() - if loop.is_running(): - loop.create_task(_rescan_sources()) - loop.create_task(_recover_analysis_jobs_on_startup()) - logger.info("Scheduler: initial source rescan scheduled immediately") + asyncio.get_running_loop() + track_background_task(_rescan_sources(), name="startup-source-rescan") + track_background_task(_recover_analysis_jobs_on_startup(), name="startup-analysis-recovery") + logger.info("Scheduler: initial source rescan scheduled immediately") except RuntimeError: logger.warning("Scheduler: could not schedule initial rescan (no event loop)") diff --git a/backend/tests/test_task_registry.py b/backend/tests/test_task_registry.py new file mode 100644 index 00000000..244730da --- /dev/null +++ b/backend/tests/test_task_registry.py @@ -0,0 +1,74 @@ +"""受管后台任务注册表回归(#72 / D-9)。 + +钉死三条契约: +- 异常任务由 done callback 记 error 日志并自动注销(不再静默丢失); +- drain 取消并等待全部受管任务,注册表清空; +- drain 超时的残留任务记 warning 且不抛出。 +""" + +from __future__ import annotations + +import asyncio +import logging + +import pytest + +from app.core import task_registry +from app.core.task_registry import drain_tracked_tasks, track_background_task + + +@pytest.fixture(autouse=True) +def _clean_registry(): + task_registry._tasks.clear() + yield + task_registry._tasks.clear() + + +@pytest.mark.asyncio +async def test_failed_task_is_logged_and_discarded(caplog): + async def boom(): + raise RuntimeError("bg exploded") + + task = track_background_task(boom(), name="boom-task") + with caplog.at_level(logging.ERROR, logger="app.core.task_registry"): + await asyncio.wait({task}) # 不直接 await:那会把任务异常重抛进用例 + + assert task not in task_registry._tasks + assert any("bg exploded" in r.getMessage() for r in caplog.records) + + +@pytest.mark.asyncio +async def test_drain_cancels_and_clears_pending_tasks(): + async def sleeper(): + await asyncio.sleep(60) + + task = track_background_task(sleeper(), name="sleeper") + await drain_tracked_tasks(timeout=1.0) + + assert task.cancelled() or task.done() + assert not task_registry._tasks + + +@pytest.mark.asyncio +async def test_drain_timeout_warns_and_does_not_raise(caplog): + async def stubborn(): + try: + await asyncio.sleep(60) + except asyncio.CancelledError: + # 吞掉取消再短暂收尾,模拟不可及时取消的任务 + await asyncio.sleep(0.05) + + task = track_background_task(stubborn(), name="stubborn") + await asyncio.sleep(0) # 让任务真正起跑:未启动即 cancel 会直接终态 cancelled + with caplog.at_level(logging.WARNING, logger="app.core.task_registry"): + await drain_tracked_tasks(timeout=0.01) + + assert any("outlived shutdown drain" in r.getMessage() for r in caplog.records) + assert not task_registry._tasks + # 给残留任务机会收尾,避免泄漏到后续用例 + await asyncio.wait([task], timeout=5) + + +@pytest.mark.asyncio +async def test_drain_empty_registry_is_noop(): + await drain_tracked_tasks() # 不应抛出