"""通用长任务 runner(BackgroundTask 跑 `jobs` 表上的异步工作;ARCH §7.4)。 T4.3「学文风走 jobs」复用此基建:`POST /style` 立即写一行 `jobs(status=queued)` 返 202 `{job_id}`,再经 FastAPI `BackgroundTasks` 登记 `run_job(...)` 跑真正的提取工作。 **独立 session 纪律**(同 `services/foreshadow_scan.run_overdue_scan` 先例 + memory/gotchas): BackgroundTask 在请求-response 发回、请求 session 关闭**之后**才跑——故 `run_job` **自建新 session**(经 `session_factory`),绝不复用请求 session。 `work: Callable[[AsyncSession], Awaitable[dict]]` 是业务逻辑缝(T4.3 部分应用「跑提取 → 写 style_fingerprint → 返回 result 摘要」)。`work` 拿到的 session 与 job 状态写同一 session → 一次 `commit()` 一并落库(业务写 + job done 原子)。 可测性:`run_job` 经可注入 `session_factory`/`repo_factory` 缝——单测直接 `await` 它, 注 fake session 工厂 + fake job repo + fake work(**不起后台线程、不连真 DB**),断言 成功路置 done、异常路置 failed。 """ from __future__ import annotations import uuid from collections.abc import Awaitable, Callable from typing import Any, Protocol import structlog from sqlalchemy.ext.asyncio import AsyncSession from ww_core.domain.job_repo import JobView, SqlJobRepo from ww_api.services.foreshadow_scan import SessionFactory log = structlog.get_logger(__name__) # 业务工作缝:拿 session 跑真正的长任务,返回写回 job.result 的摘要 dict。 JobWork = Callable[[AsyncSession], Awaitable[dict[str, Any]]] class JobLifecycleRepo(Protocol): """`run_job` 对 job repo 的**最小**依赖(仅生命周期三态写)——便于注入 fake。""" async def set_running(self, job_id: uuid.UUID) -> JobView: ... async def complete(self, job_id: uuid.UUID, result: dict[str, Any]) -> JobView: ... async def fail(self, job_id: uuid.UUID, error: str) -> JobView: ... # repo 工厂:从新 session 造 job repo。默认建 SQL 实现;测试注 fake(避免真连 DB)。 JobRepoFactory = Callable[[AsyncSession], JobLifecycleRepo] def _default_repo_factory(session: AsyncSession) -> JobLifecycleRepo: return SqlJobRepo(session) async def run_job( session_factory: SessionFactory, job_id: uuid.UUID, work: JobWork, *, request_id: str | None = None, repo_factory: JobRepoFactory = _default_repo_factory, ) -> None: """跑一个长任务:新建独立 session → set_running → await work → complete/fail → commit。 成功:`complete(job_id, result)`(status=done, progress=100, result=work 返回值)后 commit。 异常:回滚 work 的部分写 → 新 session 里 `fail(job_id, str(exc))` → commit(job 失败可见)。 任何异常都被吞(后台任务边界,不冒泡崩进程);失败置态本身再炸只记日志。 `session_factory`/`repo_factory` 是可注入缝:测试直接 await、注 fake,绝不联网/起线程。 """ try: async with session_factory() as session: repo = repo_factory(session) await repo.set_running(job_id) result = await work(session) await repo.complete(job_id, result) await session.commit() log.info("job_done", job_id=str(job_id), request_id=request_id) except Exception as exc: # noqa: BLE001 — 后台任务边界:记错误 + 置 job failed,不冒泡。 log.error("job_failed", job_id=str(job_id), request_id=request_id, error=str(exc)) await _mark_failed(session_factory, job_id, str(exc), repo_factory, request_id) async def _mark_failed( session_factory: SessionFactory, job_id: uuid.UUID, error: str, repo_factory: JobRepoFactory, request_id: str | None, ) -> None: """在一个**全新** session 里把 job 置 failed(前一 session 的事务已因异常作废)。""" try: async with session_factory() as session: repo = repo_factory(session) await repo.fail(job_id, error) await session.commit() except Exception as exc: # noqa: BLE001 — 置失败态本身再炸只记日志,不冒泡。 log.error("job_fail_mark_failed", job_id=str(job_id), request_id=request_id, error=str(exc))