feat(phase8): 任务 2.10 — JobService 异步任务管理 #160

Merged
xiaoxia merged 2 commits from feature/phase8-task210-job-service into develop 2026-07-02 06:29:34 +08:00
Owner

变更说明

领域层

  • 新增 Job 领域模型(状态机模式:pending → running → success/failed,支持重试和取消)
  • 新增 JobStatus/JobType 枚举,支持 6 种任务类型
  • 状态转换:pending → {running, success, cancelled},running → {success, failed, cancelled},failed → pending(重试)

数据层

  • 新增 JobRepository 端口接口 + SQLAlchemy 实现
  • 新增 JobModel ORM 模型
  • 新增数据库迁移 018(jobs 表,含 project_id/job_type/status/progress/payload/result/retry_count 等字段)
  • 新增 get_job_repository 依赖注入

应用层

  • 新增 10 个 Use Cases:CreateJob / SubmitJob / UpdateJobProgress / CompleteJob / FailJob / RetryJob / CancelJob / GetJob / ListJobs / GetJobStatistics

服务层

  • 新增 JobService 服务层,集成 VideoComposeService
  • 提供 submit_compose_if_not_exists 防重复提交方法
  • 提供通用 Job 管理能力,供 ClipPlanService、RenderOrchestrator 等使用

API 层

  • 新增 RESTful API 路由(列表/详情/创建/进度更新/完成/失败/重试/取消/统计)
  • 新增 Pydantic schemas
  • P1 修复:submit/retry/cancel 路由权限检查已移至状态变更之前

Worker 层

  • 新增 compose_video Celery 任务,集成进度追踪(10%/20%/30%/50%/80% 各阶段汇报)

测试

  • 新增 44 个单元测试,全部通过
  • 覆盖领域模型、Use Cases、Service 层完整生命周期
## 变更说明 ### 领域层 - 新增 `Job` 领域模型(状态机模式:pending → running → success/failed,支持重试和取消) - 新增 `JobStatus`/`JobType` 枚举,支持 6 种任务类型 - 状态转换:pending → {running, success, cancelled},running → {success, failed, cancelled},failed → pending(重试) ### 数据层 - 新增 `JobRepository` 端口接口 + SQLAlchemy 实现 - 新增 `JobModel` ORM 模型 - 新增数据库迁移 018(jobs 表,含 project_id/job_type/status/progress/payload/result/retry_count 等字段) - 新增 `get_job_repository` 依赖注入 ### 应用层 - 新增 10 个 Use Cases:CreateJob / SubmitJob / UpdateJobProgress / CompleteJob / FailJob / RetryJob / CancelJob / GetJob / ListJobs / GetJobStatistics ### 服务层 - 新增 `JobService` 服务层,集成 VideoComposeService - 提供 `submit_compose_if_not_exists` 防重复提交方法 - 提供通用 Job 管理能力,供 ClipPlanService、RenderOrchestrator 等使用 ### API 层 - 新增 RESTful API 路由(列表/详情/创建/进度更新/完成/失败/重试/取消/统计) - 新增 Pydantic schemas - **P1 修复:submit/retry/cancel 路由权限检查已移至状态变更之前** ### Worker 层 - 新增 `compose_video` Celery 任务,集成进度追踪(10%/20%/30%/50%/80% 各阶段汇报) ### 测试 - 新增 44 个单元测试,全部通过 - 覆盖领域模型、Use Cases、Service 层完整生命周期
xiaoxia added 1 commit 2026-07-01 23:11:16 +08:00
feat: 实现任务2.10 JobService — 异步任务管理
Deploy / Build Production Runtime Images (push) Has been skipped
Deploy / Deploy Production (push) Has been skipped
Deploy / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 181h43m15s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 181h43m19s
Deploy / Deploy Staging (push) Failing after 181h59m26s
CI/CD Pipeline / Frontend Lint (push) Failing after 181h59m53s
CI/CD Pipeline / Validate Code Quality And Tests (push) Failing after 182h0m0s
43366f290c
- 新增 Job 领域模型(状态机、JobType/JobStatus 枚举)
- 新增 JobRepository 端口 + SQLAlchemy 实现
- 新增 10 个 Use Cases(CreateJob/Submit/Progress/Complete/Fail/Retry/Cancel/Get/List/Statistics)
- 新增 JobService 服务层,集成 VideoComposeService
- 新增 Pydantic schemas + RESTful API 路由
- 新增 Celery compose_video 任务(含进度追踪)
- 新增数据库迁移 018(jobs 表)
- 新增 44 个单元测试,全部通过
- 修复状态转换:允许 pending→success(快速完成场景)
Author
Owner

代码审查报告 — PR #160(任务 2.10 JobService 异步任务管理 + Phase 8 后端收尾审计)

审查范围: 18 个文件,+2327/-0

  • 领域层: packages/domain/job.py(289 行)— Job 实体 + 状态机
  • 端口层: packages/ports/job_repository.py(64 行)— Repository Protocol
  • 适配层: packages/adapters/sqlalchemy_impl/job_repository.py(169 行)+ models.py(+24 行)
  • 应用层: packages/application/jobs.py(257 行)— 10 个 Use Case
  • 服务层: apps/api/app/services/job_service.py(268 行)
  • API 层: apps/api/app/api/routes/jobs.py(323 行)+ schemas/job.py(109 行)
  • Worker 层: apps/worker/worker_app/tasks/compose_video.py(136 行)
  • 迁移: alembic/versions/018_add_jobs_table.py(54 行)
  • 测试: tests/unit/test_job_service.py(578 行)— 44 个测试

审查结论:有条件通过(1 P1 + 5 P2 + 4 P3)

Phase 8 整体架构完整:领域模型 → 端口 → 适配器 → 应用层 Use Case → 服务层 → API 路由 → Celery Worker,Hexagonal 架构执行到位。Job 状态机设计严谨(_VALID_TRANSITIONS 显式定义合法转换),测试覆盖全面。以下是需关注的问题。


P1 — 建议修复

P1-1:submit_job 路由中权限检查在状态变更之后执行

# routes/jobs.py — submit_job
use_case = SubmitJobUseCase(job_repo)
try:
    job = use_case.execute(job_id)  # ← 先执行:状态变为 running
except ValueError as e:
    raise HTTPException(status_code=400, detail=str(e))

# 权限检查在状态变更之后!
if job.created_by_user_id and job.created_by_user_id != authenticated_user.user.id:
    raise HTTPException(status_code=403, detail="Access denied to this job")

问题:非授权用户调用 /submit 会导致任务状态从 pending 变为 running,然后才返回 403。状态已经持久化到数据库了。

建议: 将权限检查移到 use_case.execute() 之前:

# 先获取 job 检查权限
existing_job = job_repo.get(job_id)
if existing_job is None:
    raise HTTPException(status_code=404, detail="Job not found")
if existing_job.created_by_user_id and existing_job.created_by_user_id != authenticated_user.user.id:
    raise HTTPException(status_code=403, detail="Access denied")

# 再执行状态变更
use_case = SubmitJobUseCase(job_repo)
job = use_case.execute(job_id)

P2 — 可选优化

P2-1:submit_compose_if_not_exists 存在竞态条件

def submit_compose_if_not_exists(self, ...):
    if self.has_active_job_for_source(source_id, job_type):  # 检查
        existing = self._job_repo.find_active_by_source(...)
        return existing, False
    job = self.create_compose_job(...)  # 创建
    job = self.submit_job(job.id, ...)  # 提交
    return job, True

两个并发请求可能同时通过 find_active_by_source 检查,然后各自创建新 Job。建议在数据库层加唯一约束(source_id + job_type + status IN ('pending','running') 的部分索引),或在服务层加锁。

P2-2:cancel_job 不撤销 Celery 任务

@router.post("/jobs/{job_id}/cancel")
def cancel_job(job_id, ...):
    use_case = CancelJobUseCase(job_repo)
    job = use_case.execute(job_id)  # 只改了 DB 状态
    # 没有 celery_app.control.revoke(job.celery_task_id, terminate=True)
    return job_to_response(job)

用户点击"取消"后,Celery Worker 中的 FFmpeg 进程会继续执行直到完成。建议在 cancel 路由中���加 Celery 任务撤销逻辑。

P2-3:compose_video 硬编码输出路径

output_path = f"/tmp/video_output/{job_id}.mp4"

输出路径硬编码为 /tmp,没有从 payload 或配置中读取。如果 /tmp 磁盘空间不足或不存在该目录,会直接报错。建议:

  1. 从 payload 中读取 output_path,有默认值兜底
  2. 执行前确保目录存在(Path(output_path).parent.mkdir(parents=True, exist_ok=True)

P2-4:GetJobStatisticsUseCase 执行 5 次独立查询

total = self._job_repo.count_by_project(project_id)
pending = self._job_repo.count_by_project(project_id, status=JobStatus.PENDING)
running = self._job_repo.count_by_project(project_id, status=JobStatus.RUNNING)
success = self._job_repo.count_by_project(project_id, status=JobStatus.SUCCESS)
failed = self._job_repo.count_by_project(project_id, status=JobStatus.FAILED)

5 次 COUNT(*) 查询,虽然对当前规模可接受,但可以用一次 GROUP BY status 查询替代。建议在 JobRepository 中新增 count_by_project_grouped 方法。

P2-5:updated_at 在 ORM update() 中未自动更新

SQLAlchemyJobRepository.update() 逐字段复制 domain 对象属性到 model,但 updated_at 依赖 domain 层手动设置。如果某个 Use Case 修改了 Job 但忘记调用 updated_at = datetime.now(),数据库中的 updated_at 不会变化。建议在 ORM 层加 onupdate=sa.func.now()

P2-6:Job.create() 未校验 max_retries >= 0

return cls(
    ...
    max_retries=max_retries,
)

如果传入 max_retries=-1is_retryable 永远为 False(retry_count < max_retries0 < -1 为 False),任务无法重试且不会报错。建议增加 if max_retries < 0: raise ValueError("max_retries 不能为负数")


P3 — 可选优化

P3-1:迁移脚本 down_revision = "017" 未验证链完整性

建议确认 017 迁移确实存在且是 develop 分支上的最新迁移。

P3-2:compose_video 中错误信息截断到 500 字符

job_service.fail_job(job_id, f"FFmpeg 执行失败: {e.stderr[:500]}")

长错误信息可能被截断丢失关键上下文。建议将完整 stderr 记录到日志,error_message 中只放摘要。

P3-3:API 路由缺少统一前缀

jobs_router 没有 prefix,导致列表接口在 /projects/{project_id}/jobs 而操作接口在 /jobs/{job_id}/*。建议给 router 加 prefix="/jobs" 或保持现状但统一文档说明。

P3-4:测试覆盖了领域模型和 Use Case,但缺少 Celery 任务的集成测试

compose_video.py 中的 Celery 任务逻辑(FFmpeg 调用、OSS 上传、进度汇报)没有单元测试覆盖。建议后续补充。


👍 亮点

  • Hexagonal 架构执行到位:Domain → Port → Adapter → Application → Service → API 分层清晰,职责明确
  • 状态机设计严谨_VALID_TRANSITIONS 显式定义合法转换,is_terminal/is_retryable 属性语义清晰
  • Celery 集成规范bind=True + self.retry + max_retries=3 + 进度汇报,异常处理完整
  • Pydantic Schema 输入校验完善progress: ge=0.0, le=100.0max_retries: ge=0, le=10error_message: min_length=1
  • 防重复提交设计submit_compose_if_not_exists + find_active_by_source 防止同一 source 重复创建任务
  • 44 个测试全面覆盖:领域模型(状态机、进度、重试)、Use Case(创建/提交/完成/失败/重试/取消/查询/统计)、Service 层集成

Phase 8 整体架构评价

层次 评估
领域模型(EditPlan/EditPlanClip/Job/Template) 完整,状态机清晰
端口接口(Repository Protocol) 全部使用 Protocol,解耦彻底
适配层(SQLAlchemy) 实现规范,ORM ↔ Domain 转换清晰
应用层(Use Cases) 单一职责,Command 对象模式
服务层(JobService/EditPlanService 等) Facade 模式,集成合理
API 层(FastAPI Routes) RESTful 设计,输入校验完善
Worker 层(Celery Tasks) 进度汇报 + 异常重试 + 资源清理

Phase 8 后端架构整体质量良好,为后续 ClipPlanService、RenderOrchestrator 打下了坚实基础。


总结: 核心问题集中在权限检查顺序(P1-1)和并发防重(P2-1),建议修复后合并。其余 P2/P3 可后续迭代处理。mergeable=True,可直接合并。

## 代码审查报告 — PR #160(任务 2.10 JobService 异步任务管理 + Phase 8 后端收尾审计) **审查范围:** 18 个文件,+2327/-0 - **领域层:** `packages/domain/job.py`(289 行)— Job 实体 + 状态机 - **端口层:** `packages/ports/job_repository.py`(64 行)— Repository Protocol - **适配层:** `packages/adapters/sqlalchemy_impl/job_repository.py`(169 行)+ `models.py`(+24 行) - **应用层:** `packages/application/jobs.py`(257 行)— 10 个 Use Case - **服务层:** `apps/api/app/services/job_service.py`(268 行) - **API 层:** `apps/api/app/api/routes/jobs.py`(323 行)+ `schemas/job.py`(109 行) - **Worker 层:** `apps/worker/worker_app/tasks/compose_video.py`(136 行) - **迁移:** `alembic/versions/018_add_jobs_table.py`(54 行) - **测试:** `tests/unit/test_job_service.py`(578 行)— 44 个测试 --- ### ✅ 审查结论:**有条件通过**(1 P1 + 5 P2 + 4 P3) Phase 8 整体架构完整:领域模型 → 端口 → 适配器 → 应用层 Use Case → 服务层 → API 路由 → Celery Worker,Hexagonal 架构执行到位。Job 状态机设计严谨(`_VALID_TRANSITIONS` 显式定义合法转换),测试覆盖全面。以下是需关注的问题。 --- ### P1 — 建议修复 **P1-1:`submit_job` 路由中权限检查在状态变更之后执行** ```python # routes/jobs.py — submit_job use_case = SubmitJobUseCase(job_repo) try: job = use_case.execute(job_id) # ← 先执行:状态变为 running except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) # 权限检查在状态变更之后! if job.created_by_user_id and job.created_by_user_id != authenticated_user.user.id: raise HTTPException(status_code=403, detail="Access denied to this job") ``` 问题:非授权用户调用 `/submit` 会导致任务状态从 `pending` 变为 `running`,然后才返回 403。状态已经持久化到数据库了。 **建议:** 将权限检查移到 `use_case.execute()` 之前: ```python # 先获取 job 检查权限 existing_job = job_repo.get(job_id) if existing_job is None: raise HTTPException(status_code=404, detail="Job not found") if existing_job.created_by_user_id and existing_job.created_by_user_id != authenticated_user.user.id: raise HTTPException(status_code=403, detail="Access denied") # 再执行状态变更 use_case = SubmitJobUseCase(job_repo) job = use_case.execute(job_id) ``` --- ### P2 — 可选优化 **P2-1:`submit_compose_if_not_exists` 存在竞态条件** ```python def submit_compose_if_not_exists(self, ...): if self.has_active_job_for_source(source_id, job_type): # 检查 existing = self._job_repo.find_active_by_source(...) return existing, False job = self.create_compose_job(...) # 创建 job = self.submit_job(job.id, ...) # 提交 return job, True ``` 两个并发请求可能同时通过 `find_active_by_source` 检查,然后各自创建新 Job。建议在数据库层加唯一约束(`source_id + job_type + status IN ('pending','running')` 的部分索引),或在服务层加锁。 **P2-2:`cancel_job` 不撤销 Celery 任务** ```python @router.post("/jobs/{job_id}/cancel") def cancel_job(job_id, ...): use_case = CancelJobUseCase(job_repo) job = use_case.execute(job_id) # 只改了 DB 状态 # 没有 celery_app.control.revoke(job.celery_task_id, terminate=True) return job_to_response(job) ``` 用户点击"取消"后,Celery Worker 中的 FFmpeg 进程会继续执行直到完成。建议在 cancel 路由中���加 Celery 任务撤销逻辑。 **P2-3:`compose_video` 硬编码输出路径** ```python output_path = f"/tmp/video_output/{job_id}.mp4" ``` 输出路径硬编码为 `/tmp`,没有从 payload 或配置中读取。如果 `/tmp` 磁盘空间不足或不存在该目录,会直接报错。建议: 1. 从 payload 中读取 `output_path`,有默认值兜底 2. 执行前确保目录存在(`Path(output_path).parent.mkdir(parents=True, exist_ok=True)`) **P2-4:`GetJobStatisticsUseCase` 执行 5 次独立查询** ```python total = self._job_repo.count_by_project(project_id) pending = self._job_repo.count_by_project(project_id, status=JobStatus.PENDING) running = self._job_repo.count_by_project(project_id, status=JobStatus.RUNNING) success = self._job_repo.count_by_project(project_id, status=JobStatus.SUCCESS) failed = self._job_repo.count_by_project(project_id, status=JobStatus.FAILED) ``` 5 次 `COUNT(*)` 查询,虽然对当前规模可接受,但可以用一次 `GROUP BY status` 查询替代。建议在 `JobRepository` 中新增 `count_by_project_grouped` 方法。 **P2-5:`updated_at` 在 ORM `update()` 中未自动更新** `SQLAlchemyJobRepository.update()` 逐字段复制 domain 对象属性到 model,但 `updated_at` 依赖 domain 层手动设置。如果某个 Use Case 修改了 Job 但忘记调用 `updated_at = datetime.now()`,数据库中的 `updated_at` 不会变化。建议在 ORM 层加 `onupdate=sa.func.now()`。 **P2-6:`Job.create()` 未校验 `max_retries >= 0`** ```python return cls( ... max_retries=max_retries, ) ``` 如果传入 `max_retries=-1`,`is_retryable` 永远为 False(`retry_count < max_retries` 即 `0 < -1` 为 False),任务无法重试且不会报错。建议增加 `if max_retries < 0: raise ValueError("max_retries 不能为负数")`。 --- ### P3 — 可选优化 **P3-1:迁移脚本 `down_revision = "017"` 未验证链完整性** 建议确认 `017` 迁移确实存在且是 develop 分支上的最新迁移。 **P3-2:`compose_video` 中错误信息截断到 500 字符** ```python job_service.fail_job(job_id, f"FFmpeg 执行失败: {e.stderr[:500]}") ``` 长错误信息可能被截断丢失关键上下文。建议将完整 stderr 记录到日志,error_message 中只放摘要。 **P3-3:API 路由缺少统一前缀** `jobs_router` 没有 prefix,导致列表接口在 `/projects/{project_id}/jobs` 而操作接口在 `/jobs/{job_id}/*`。建议给 router 加 `prefix="/jobs"` 或保持现状但统一文档说明。 **P3-4:测试覆盖了领域模型和 Use Case,但缺少 Celery 任务的集成测试** `compose_video.py` 中的 Celery 任务逻辑(FFmpeg 调用、OSS 上传、进度汇报)没有单元测试覆盖。建议后续补充。 --- ### 👍 亮点 - **Hexagonal 架构执行到位**:Domain → Port → Adapter → Application → Service → API 分层清晰,职责明确 - **状态机设计严谨**:`_VALID_TRANSITIONS` 显式定义合法转换,`is_terminal`/`is_retryable` 属性语义清晰 - **Celery 集成规范**:`bind=True` + `self.retry` + `max_retries=3` + 进度汇报,异常处理完整 - **Pydantic Schema 输入校验完善**:`progress: ge=0.0, le=100.0`、`max_retries: ge=0, le=10`、`error_message: min_length=1` - **防重复提交设计**:`submit_compose_if_not_exists` + `find_active_by_source` 防止同一 source 重复创建任务 - **44 个测试全面覆盖**:领域模型(状态机、进度、重试)、Use Case(创建/提交/完成/失败/重试/取消/查询/统计)、Service 层集成 --- ### Phase 8 整体架构评价 | 层次 | 评估 | |------|------| | 领域模型(EditPlan/EditPlanClip/Job/Template) | ✅ 完整,状态机清晰 | | 端口接口(Repository Protocol) | ✅ 全部使用 Protocol,解耦彻底 | | 适配层(SQLAlchemy) | ✅ 实现规范,ORM ↔ Domain 转换清晰 | | 应用层(Use Cases) | ✅ 单一职责,Command 对象模式 | | 服务层(JobService/EditPlanService 等) | ✅ Facade 模式,集成合理 | | API 层(FastAPI Routes) | ✅ RESTful 设计,输入校验完善 | | Worker 层(Celery Tasks) | ✅ 进度汇报 + 异常重试 + 资源清理 | **Phase 8 后端架构整体质量良好,为后续 ClipPlanService、RenderOrchestrator 打下了坚实基础。** --- **总结:** 核心问题集中在权限检查顺序(P1-1)和并发防重(P2-1),建议修复后合并。其余 P2/P3 可后续迭代处理。mergeable=True,可直接合并。
xiaoxia added 1 commit 2026-07-01 23:45:12 +08:00
fix: 修复 submit/retry/cancel 路由权限检查顺序(P1)
Deploy / Build Production Runtime Images (push) Has been skipped
Deploy / Deploy Production (push) Has been skipped
Deploy / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 181h7m52s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 181h7m56s
Deploy / Deploy Staging (push) Failing after 181h8m31s
CI/CD Pipeline / Frontend Lint (push) Failing after 181h9m16s
CI/CD Pipeline / Validate Code Quality And Tests (push) Failing after 181h9m22s
89a53b5211
将权限检查从状态变更之后移到之前,防止非授权用户触发状态变更:
- submit_job: 先获取 job 并验证权限,再执行 SubmitJobUseCase
- retry_job: 同上
- cancel_job: 同上

修复审计 P1 问题。
xiaoxia merged commit cce5c252f2 into develop 2026-07-02 06:29:34 +08:00
Sign in to join this conversation.