Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions backend/app/api/v1/daily_reports.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
53 changes: 53 additions & 0 deletions backend/app/core/task_registry.py
Original file line number Diff line number Diff line change
@@ -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()
6 changes: 6 additions & 0 deletions backend/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 7 additions & 6 deletions backend/app/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)")

Expand Down
74 changes: 74 additions & 0 deletions backend/tests/test_task_registry.py
Original file line number Diff line number Diff line change
@@ -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() # 不应抛出
Loading