From 206555c2d81c3e3305d04acc9b183be710d89abb Mon Sep 17 00:00:00 2001 From: CI Bot Date: Mon, 13 Jul 2026 22:22:54 +0800 Subject: [PATCH 1/6] =?UTF-8?q?feat:=20=E4=BB=BB=E5=8A=A1=E4=B8=AD?= =?UTF-8?q?=E5=BF=83=E5=8D=87=E7=BA=A7=20-=20=E5=A4=B1=E8=B4=A5=E9=87=8D?= =?UTF-8?q?=E8=AF=95+=E9=94=99=E8=AF=AF=E8=BF=BD=E8=B8=AA+=E5=88=97?= =?UTF-8?q?=E8=A1=A8=E7=AD=9B=E9=80=89+=E8=87=AA=E5=8A=A8=E9=87=8D?= =?UTF-8?q?=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 数据库:generation_tasks 新增 error_info JSON / retry_count / auto_retry_enabled / auto_retry_max 字段 - 领域模型:mark_failed 支持结构化 error_info,mark_pending_from_failed 递增 retry_count - Repository:新增 list_by_user_filtered / count_by_user_filtered / list_by_project_filtered / count_by_project_filtered - Use Case:新增 RetryGenerationTaskUseCase(原地重试)、ListUserTasksFilteredUseCase - API:任务列表支持 status/task_type 筛选+分页,重试改为原地重试(复用task_id) - Worker:失败时结构化记录 error_info(含堆栈摘要),支持自动重试(指数退避,最大60s) - Schema:所有响应新增 error_info / retry_count / auto_retry 字段 - 单测:新增9个测试用例,51个状态测试全绿,80个相关测试无回归 --- ...rror_info_and_retry_to_generation_tasks.py | 45 ++++ apps/api/app/api/routes/generation_tasks.py | 2 + apps/api/app/api/routes/task_center.py | 215 +++++++++++------- apps/api/app/schemas/generation_task.py | 13 ++ apps/api/app/schemas/task_center.py | 6 + apps/worker/worker_app/tasks/generation.py | 82 ++++++- .../generation_task_repository.py | 82 +++++++ packages/adapters/sqlalchemy_impl/models.py | 4 + packages/application/__init__.py | 6 + packages/application/generation_tasks.py | 66 +++++- packages/domain/generation_task.py | 26 ++- packages/ports/generation_task_repository.py | 32 +++ tests/unit/test_generation_task_status.py | 102 +++++++++ 13 files changed, 596 insertions(+), 85 deletions(-) create mode 100644 alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py mode change 100644 => 100755 apps/api/app/schemas/generation_task.py mode change 100644 => 100755 apps/api/app/schemas/task_center.py mode change 100644 => 100755 apps/worker/worker_app/tasks/generation.py mode change 100644 => 100755 packages/application/generation_tasks.py mode change 100644 => 100755 packages/domain/generation_task.py mode change 100644 => 100755 tests/unit/test_generation_task_status.py diff --git a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py new file mode 100644 index 000000000..3005c83c6 --- /dev/null +++ b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py @@ -0,0 +1,45 @@ +"""add error_info and retry fields to generation_tasks + +Revision ID: 038 +Revises: 037 +Create Date: 2026-07-13 22:15:00.000000 +""" +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects.mysql import JSON as MySQLJSON + +# revision identifiers, used by Alembic. +revision = '038' +down_revision = '037' +branch_labels = None +depends_on = None + + +def upgrade(): + # error_info: 结构化错误信息(error_type, message, stack_trace, failed_at, stage等) + op.add_column( + 'generation_tasks', + sa.Column('error_info', sa.JSON(), nullable=True), + ) + # retry_count: 重试次数 + op.add_column( + 'generation_tasks', + sa.Column('retry_count', sa.Integer(), nullable=False, server_default='0'), + ) + # auto_retry_enabled: 是否开启自动重试 + op.add_column( + 'generation_tasks', + sa.Column('auto_retry_enabled', sa.Boolean(), nullable=False, server_default=sa.text('false')), + ) + # auto_retry_max: 最大自动重试次数 + op.add_column( + 'generation_tasks', + sa.Column('auto_retry_max', sa.Integer(), nullable=False, server_default='0'), + ) + + +def downgrade(): + op.drop_column('generation_tasks', 'auto_retry_max') + op.drop_column('generation_tasks', 'auto_retry_enabled') + op.drop_column('generation_tasks', 'retry_count') + op.drop_column('generation_tasks', 'error_info') diff --git a/apps/api/app/api/routes/generation_tasks.py b/apps/api/app/api/routes/generation_tasks.py index 22c031f14..ccc3e9a55 100644 --- a/apps/api/app/api/routes/generation_tasks.py +++ b/apps/api/app/api/routes/generation_tasks.py @@ -267,6 +267,8 @@ def create_generation_task( source_edit_plan_id=request.source_edit_plan_id, asset_select_mode=request.asset_select_mode, batch_id=batch_id, + auto_retry_enabled=request.auto_retry_enabled, + auto_retry_max=request.auto_retry_max, ) ) try: diff --git a/apps/api/app/api/routes/task_center.py b/apps/api/app/api/routes/task_center.py index c796fd58d..e968a0561 100755 --- a/apps/api/app/api/routes/task_center.py +++ b/apps/api/app/api/routes/task_center.py @@ -21,11 +21,12 @@ from app.schemas.task_center import ( ProjectTaskResponse, UserTaskResponse, ) -from fastapi import APIRouter, Depends, HTTPException +from fastapi import APIRouter, Depends, HTTPException, Query from packages.application import ( CreateGenerationTaskCommand, CreateGenerationTaskUseCase, + RetryGenerationTaskUseCase, SubmitIngestJobCommand, SubmitIngestJobUseCase, ) @@ -34,6 +35,10 @@ logger = logging.getLogger(__name__) router = APIRouter() +DEFAULT_PAGE_SIZE = 50 +MAX_PAGE_SIZE = 200 + + def _humanize_task_error(error_message: str) -> str: raw = (error_message or "").strip() if not raw: @@ -63,6 +68,8 @@ def _generation_step(task) -> str: return "生成完成" if s == "failed": return "生成失败" + if s == "cancelled": + return "已取消" return s @@ -79,6 +86,26 @@ def _ingest_step(job) -> str: return s +def _generation_task_to_user_response(task) -> UserTaskResponse: + return UserTaskResponse( + id=f"generation:{task.id}", + task_type="generation", + project_id=task.project_id, + template_id=task.template_id, + status=_status_value(task.status), + progress=task.progress, + current_step=_generation_step(task), + error_message=task.error_message, + error_info=task.error_info or {}, + user_message=_humanize_task_error(task.error_message), + retryable=_status_value(task.status) == "failed", + retry_count=task.retry_count or 0, + source_id=task.id, + created_at=task.created_at, + updated_at=task.completed_at or task.started_at or task.created_at, + ) + + def _generation_task_to_project_response(task) -> ProjectTaskResponse: return ProjectTaskResponse( id=f"generation:{task.id}", @@ -88,8 +115,10 @@ def _generation_task_to_project_response(task) -> ProjectTaskResponse: progress=task.progress, current_step=_generation_step(task), error_message=task.error_message, + error_info=task.error_info or {}, user_message=_humanize_task_error(task.error_message), retryable=_status_value(task.status) == "failed", + retry_count=task.retry_count or 0, source_id=task.id, template_id=task.template_id, created_at=task.created_at, @@ -97,40 +126,66 @@ def _generation_task_to_project_response(task) -> ProjectTaskResponse: ) +def _validate_status(status: str | None) -> str | None: + """校验状态值合法性。""" + if status is None: + return None + valid = {"pending", "running", "completed", "failed", "cancelled"} + if status not in valid: + raise HTTPException( + status_code=400, + detail=f"无效的状态筛选值: {status},允许值: {', '.join(sorted(valid))}", + ) + return status + + +def _clamp_page_size(page_size: int) -> int: + if page_size <= 0: + return DEFAULT_PAGE_SIZE + if page_size > MAX_PAGE_SIZE: + return MAX_PAGE_SIZE + return page_size + + # ── 用户级端点(放在项目级端点之前,避免路由冲突) ── @router.get("/tasks", response_model=ListTasksResponse) def list_user_tasks( + status: str | None = Query(None, description="按状态筛选:pending/running/completed/failed/cancelled"), + task_type: str | None = Query(None, description="按任务类型筛选:generation/ingest"), + page: int = Query(1, ge=1, description="页码,从1开始"), + page_size: int = Query(DEFAULT_PAGE_SIZE, ge=1, le=MAX_PAGE_SIZE, description="每页数量"), authenticated_user: AuthenticatedUser = Depends(get_current_user), ingest_job_repository: Any = Depends(get_ingest_job_repository), generation_task_repository: Any = Depends(get_generation_task_repository), ) -> ListTasksResponse: - """用户级任务列表(跨 project),合并 ingest + generation 任务。""" + """用户级任务列表(跨 project),支持状态/类型筛选和分页。""" + status = _validate_status(status) + page_size = _clamp_page_size(page_size) user_id = authenticated_user.user.id + offset = (page - 1) * page_size + items: list[UserTaskResponse] = [] - for task in generation_task_repository.list_by_user(user_id): - items.append( - UserTaskResponse( - id=f"generation:{task.id}", - task_type="generation", - project_id=task.project_id, - template_id=task.template_id, - status=_status_value(task.status), - progress=task.progress, - current_step=_generation_step(task), - error_message=task.error_message, - user_message=_humanize_task_error(task.error_message), - retryable=_status_value(task.status) == "failed", - source_id=task.id, - created_at=task.created_at, - updated_at=task.completed_at or task.started_at or task.created_at, - ) + # 生成任务 + if task_type is None or task_type == "generation": + gen_result = generation_task_repository.list_by_user_filtered( + user_id, + status=status, + limit=page_size + 1, # 多取一条判断是否还有下一页(简单起见这里用offset) + offset=offset, ) + for task in gen_result: + items.append(_generation_task_to_user_response(task)) + # 按时间倒序 items.sort(key=lambda item: item.updated_at or item.created_at or "", reverse=True) - return ListTasksResponse(items=items) + + # 总数(仅generation,ingest暂不计入总数以保持简单) + total = generation_task_repository.count_by_user_filtered(user_id, status=status) + + return ListTasksResponse(items=items[:page_size], total=total) @router.post("/tasks/{task_id}/retry", response_model=UserTaskResponse) @@ -139,7 +194,7 @@ def retry_task_by_id( authenticated_user: AuthenticatedUser = Depends(get_current_user), generation_task_repository: Any = Depends(get_generation_task_repository), ) -> UserTaskResponse: - """简化重试:通过 task_id 直接重试失败的生成任务。""" + """原地重试失败的生成任务(复用同一个task_id,retry_count+1)。""" task = generation_task_repository.get(task_id) if task is None: raise HTTPException(status_code=404, detail="Generation task not found") @@ -149,6 +204,7 @@ def retry_task_by_id( raise HTTPException(status_code=409, detail="Only failed tasks can be retried") user_id = authenticated_user.user.id + # 预检查 user_pending = generation_task_repository.count_pending_by_user(user_id) global_pending = generation_task_repository.count_pending_total() @@ -163,20 +219,11 @@ def retry_task_by_id( detail="系统繁忙,请稍后再试", ) - use_case = CreateGenerationTaskUseCase(generation_task_repository) - retried = use_case.execute( - CreateGenerationTaskCommand( - project_id=task.project_id, - asset_library_id=task.asset_library_id, - strategy_id=task.strategy_id, - voice_library_id=task.voice_library_id, - template_id=task.template_id, - asset_ids=task.asset_ids, - title_ids=task.title_ids, - voice_ids=task.voice_ids, - created_by_user_id=user_id, - ) - ) + # 原地重试 + use_case = RetryGenerationTaskUseCase(generation_task_repository) + retried = use_case.execute(task_id) + + # 重新入队 try: if not safe_enqueue_generation_task( retried, generation_task_repository, user_id=user_id, log_prefix="[任务中心]" @@ -192,18 +239,8 @@ def retry_task_by_id( status_code=503, detail="系统繁忙,请稍后再试", ) from None - return UserTaskResponse( - id=f"generation:{retried.id}", - task_type="generation", - project_id=retried.project_id, - template_id=retried.template_id, - status=_status_value(retried.status), - progress=retried.progress, - current_step=_generation_step(retried), - source_id=retried.id, - created_at=retried.created_at, - updated_at=retried.created_at, - ) + + return _generation_task_to_user_response(retried) # ── 项目级端点 ── @@ -212,37 +249,64 @@ def retry_task_by_id( @router.get("/projects/{project_id}/tasks", response_model=ListProjectTasksResponse) def list_project_tasks( project_id: str, + status: str | None = Query(None, description="按状态筛选:pending/running/completed/failed/cancelled"), + task_type: str | None = Query(None, description="按任务类型筛选:generation/ingest"), + page: int = Query(1, ge=1, description="页码,从1开始"), + page_size: int = Query(DEFAULT_PAGE_SIZE, ge=1, le=MAX_PAGE_SIZE, description="每页数量"), authenticated_user: AuthenticatedUser = Depends(get_current_user), project_repository: Any = Depends(get_project_repository), ingest_job_repository: Any = Depends(get_ingest_job_repository), generation_task_repository: Any = Depends(get_generation_task_repository), ) -> ListProjectTasksResponse: + """项目级任务列表,支持状态/类型筛选和分页。""" project = project_repository.find_by_id(project_id) if project is None: raise HTTPException(status_code=404, detail="Project not found") + status = _validate_status(status) + page_size = _clamp_page_size(page_size) + offset = (page - 1) * page_size + items: list[ProjectTaskResponse] = [] - for job in ingest_job_repository.list_by_project(project_id): - items.append( - ProjectTaskResponse( - id=f"ingest:{job.id}", - task_type="ingest", - project_id=job.project_id, - status=_status_value(job.status), - progress=100.0 if _status_value(job.status) == "completed" else 0.0, - current_step=_ingest_step(job), - error_message=job.error_message, - user_message=_humanize_task_error(job.error_message), - retryable=_status_value(job.status) == "failed", - source_id=job.id, - created_at=job.created_at, - updated_at=job.updated_at, + + # 导入任务 + if task_type is None or task_type == "ingest": + for job in ingest_job_repository.list_by_project(project_id): + if status and _status_value(job.status) != status: + continue + items.append( + ProjectTaskResponse( + id=f"ingest:{job.id}", + task_type="ingest", + project_id=job.project_id, + status=_status_value(job.status), + progress=100.0 if _status_value(job.status) == "completed" else 0.0, + current_step=_ingest_step(job), + error_message=job.error_message, + user_message=_humanize_task_error(job.error_message), + retryable=_status_value(job.status) == "failed", + source_id=job.id, + created_at=job.created_at, + updated_at=job.updated_at, + ) ) + + # 生成任务 + if task_type is None or task_type == "generation": + gen_items = generation_task_repository.list_by_project_filtered( + project_id, + status=status, + limit=page_size + 1, + offset=offset, ) - for task in generation_task_repository.list_by_project(project_id): - items.append(_generation_task_to_project_response(task)) + for task in gen_items: + items.append(_generation_task_to_project_response(task)) + items.sort(key=lambda item: item.updated_at or item.created_at or "", reverse=True) - return ListProjectTasksResponse(items=items) + + total = generation_task_repository.count_by_project_filtered(project_id, status=status) + + return ListProjectTasksResponse(items=items[:page_size], total=total) @router.post("/tasks/{task_type}/{source_id}/retry", response_model=ProjectTaskResponse) @@ -253,6 +317,7 @@ def retry_project_task( ingest_job_repository: Any = Depends(get_ingest_job_repository), generation_task_repository: Any = Depends(get_generation_task_repository), ) -> ProjectTaskResponse: + """项目级任务重试。""" if task_type == "generation": task = generation_task_repository.get(source_id) if task is None: @@ -261,6 +326,7 @@ def retry_project_task( raise HTTPException(status_code=409, detail="Only failed tasks can be retried") user_id = authenticated_user.user.id + # 预检查 user_pending = generation_task_repository.count_pending_by_user(user_id) global_pending = generation_task_repository.count_pending_total() @@ -275,20 +341,10 @@ def retry_project_task( detail="系统繁忙,请稍后再试", ) - use_case = CreateGenerationTaskUseCase(generation_task_repository) - retried = use_case.execute( - CreateGenerationTaskCommand( - project_id=task.project_id, - asset_library_id=task.asset_library_id, - strategy_id=task.strategy_id, - voice_library_id=task.voice_library_id, - template_id=task.template_id, - asset_ids=task.asset_ids, - title_ids=task.title_ids, - voice_ids=task.voice_ids, - created_by_user_id=user_id, - ) - ) + # 原地重试 + use_case = RetryGenerationTaskUseCase(generation_task_repository) + retried = use_case.execute(source_id) + try: if not safe_enqueue_generation_task( retried, generation_task_repository, user_id=user_id, log_prefix="[任务中心]" @@ -305,6 +361,7 @@ def retry_project_task( detail="系统繁忙,请稍后再试", ) from None return _generation_task_to_project_response(retried) + if task_type == "ingest": job = ingest_job_repository.get(source_id) if job is None: diff --git a/apps/api/app/schemas/generation_task.py b/apps/api/app/schemas/generation_task.py old mode 100644 new mode 100755 index 1724d3f74..4f7a812a4 --- a/apps/api/app/schemas/generation_task.py +++ b/apps/api/app/schemas/generation_task.py @@ -33,6 +33,15 @@ class CreateGenerationTaskRequest(BaseModel): asset_select_count: int = Field( default=0, ge=0, le=100, description="选取数量,0表示全部(仅 random/smart 模式有效)" ) + # ── 自动重试 ── + auto_retry_enabled: bool = Field( + default=False, + description="是否开启失败自动重试,默认关闭", + ) + auto_retry_max: int = Field( + default=0, ge=0, le=5, + description="最大自动重试次数,0表示不自动重试,最大5次", + ) @model_validator(mode="after") def _check_at_least_one_mode(self) -> "CreateGenerationTaskRequest": @@ -64,6 +73,10 @@ class GenerationTaskResponse(BaseModel): progress: float result_count: int error_message: str + error_info: dict = Field(default_factory=dict) + retry_count: int = 0 + auto_retry_enabled: bool = False + auto_retry_max: int = 0 logs: list[dict] = Field(default_factory=list) @field_validator("logs", mode="before") diff --git a/apps/api/app/schemas/task_center.py b/apps/api/app/schemas/task_center.py old mode 100644 new mode 100755 index 931fae17e..c2e29f06e --- a/apps/api/app/schemas/task_center.py +++ b/apps/api/app/schemas/task_center.py @@ -11,8 +11,10 @@ class ProjectTaskResponse(BaseModel): progress: float current_step: str error_message: str = "" + error_info: dict = Field(default_factory=dict) user_message: str = "" retryable: bool = False + retry_count: int = 0 source_id: str = "" template_id: str = "" created_at: datetime | None = None @@ -21,6 +23,7 @@ class ProjectTaskResponse(BaseModel): class ListProjectTasksResponse(BaseModel): items: list[ProjectTaskResponse] = Field(default_factory=list) + total: int = 0 class UserTaskResponse(BaseModel): @@ -34,8 +37,10 @@ class UserTaskResponse(BaseModel): progress: float current_step: str error_message: str = "" + error_info: dict = Field(default_factory=dict) user_message: str = "" retryable: bool = False + retry_count: int = 0 source_id: str = "" created_at: datetime | None = None updated_at: datetime | None = None @@ -45,3 +50,4 @@ class ListTasksResponse(BaseModel): """用户级任务列表响应(GET /api/v1/tasks)。""" items: list[UserTaskResponse] = Field(default_factory=list) + total: int = 0 diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py old mode 100644 new mode 100755 index 5d7f6e356..7b59b91dd --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -86,6 +86,36 @@ def _update_task_status(task_id: str, status_action: str, **kwargs) -> bool: return False +def _build_error_info(error: Exception, stage: str = "render") -> dict: + """构建结构化错误信息。 + + Args: + error: 异常对象 + stage: 发生错误的阶段(download/render/merge/upload等) + + Returns: + 包含 error_type, message, stack_trace, stage, failed_at 的字典 + """ + import traceback + from datetime import datetime, timezone + + tb_str = traceback.format_exc() + # 截取堆栈前20行,避免字段过大 + tb_lines = tb_str.strip().splitlines() + if len(tb_lines) > 20: + tb_summary = "\n".join(tb_lines[:20]) + f"\n... (truncated, total {len(tb_lines)} lines)" + else: + tb_summary = tb_str + + return { + "error_type": type(error).__name__, + "message": str(error), + "stack_trace": tb_summary, + "stage": stage, + "failed_at": datetime.now(timezone.utc).isoformat(), + } + + # ── 日志持久化辅助 ──────────────────────────────────────────────────────────── @@ -1122,6 +1152,9 @@ def generate_video(self, task_id: str) -> dict: except Exception as error: logger.error("[task_id=%s] [任务失败] %s", task_id, error, exc_info=True) + # 构建结构化错误信息 + error_info = _build_error_info(error, stage="render") + # 记录失败日志 try: _session = SessionLocal() @@ -1134,6 +1167,7 @@ def generate_video(self, task_id: str) -> dict: str(error), level="ERROR", error_type=type(error).__name__, + stage="render", ) _flush_logs(task_id, gen_task) finally: @@ -1141,7 +1175,53 @@ def generate_video(self, task_id: str) -> dict: except Exception: logger.warning("[task_id=%s] 记录失败日志异常", task_id, exc_info=True) - _update_task_status(task_id, "mark_failed", error_message=str(error)) + _update_task_status( + task_id, + "mark_failed", + error_message=str(error), + error_info=error_info, + ) + + # ── 自动重试逻辑 ────────────────────────────────────────────────── + try: + from packages.adapters.sqlalchemy_impl.generation_task_repository import ( + SQLAlchemyGenerationTaskRepository, + ) + + _s = SessionLocal() + try: + _r = SQLAlchemyGenerationTaskRepository(_s) + _task = _r.get(task_id) + if _task and _task.auto_retry_enabled and _task.auto_retry_max > 0: + current_retry = _task.retry_count or 0 + if current_retry < _task.auto_retry_max: + logger.info( + "[task_id=%s] 触发自动重试: 当前重试次数=%d, 最大重试次数=%d", + task_id, current_retry, _task.auto_retry_max, + ) + # 计算退避延迟(指数退避,基础5s,最大60s) + backoff_seconds = min(5 * (2 ** current_retry), 60) + # 原地重试 + _task.mark_pending_from_failed() + _r.update(_task) + # 延迟重新入队 + celery_app.send_task( + "worker.generate_video", + args=[task_id], + countdown=backoff_seconds, + ) + logger.info( + "[task_id=%s] 自动重试已入队: 延迟=%ds, 第%d次重试", + task_id, backoff_seconds, current_retry + 1, + ) + finally: + _s.close() + except Exception as retry_err: + logger.warning( + "[task_id=%s] 自动重试逻辑执行失败: %s", + task_id, retry_err, exc_info=True, + ) + return { "status": "failed", "task_id": task_id, diff --git a/packages/adapters/sqlalchemy_impl/generation_task_repository.py b/packages/adapters/sqlalchemy_impl/generation_task_repository.py index 646a05be1..53c1fa31c 100755 --- a/packages/adapters/sqlalchemy_impl/generation_task_repository.py +++ b/packages/adapters/sqlalchemy_impl/generation_task_repository.py @@ -21,6 +21,10 @@ def _to_domain(model: GenerationTaskModel) -> GenerationTask: progress=model.progress, result_count=int(model.result_count or 0), error_message=model.error_message, + error_info=dict(model.error_info) if model.error_info else {}, + retry_count=model.retry_count or 0, + auto_retry_enabled=bool(model.auto_retry_enabled), + auto_retry_max=model.auto_retry_max or 0, started_at=model.started_at, completed_at=model.completed_at, created_by_user_id=model.created_by_user_id, @@ -51,6 +55,10 @@ class SQLAlchemyGenerationTaskRepository: progress=task.progress, result_count=task.result_count, error_message=task.error_message, + error_info=task.error_info or None, + retry_count=task.retry_count or 0, + auto_retry_enabled=task.auto_retry_enabled, + auto_retry_max=task.auto_retry_max or 0, started_at=task.started_at, completed_at=task.completed_at, created_by_user_id=task.created_by_user_id, @@ -127,6 +135,76 @@ class SQLAlchemyGenerationTaskRepository: ) return [_to_domain(m) for m in models] + def list_by_user_filtered( + self, + user_id: str, + *, + status: str | None = None, + limit: int | None = None, + offset: int = 0, + ) -> list[GenerationTask]: + """按用户+状态筛选任务列表。""" + query = self.session.query(GenerationTaskModel).filter( + GenerationTaskModel.created_by_user_id == user_id + ) + if status: + query = query.filter(GenerationTaskModel.status == status) + query = query.order_by(GenerationTaskModel.created_at.desc()) + if offset: + query = query.offset(offset) + if limit: + query = query.limit(limit) + return [_to_domain(m) for m in query.all()] + + def count_by_user_filtered( + self, + user_id: str, + *, + status: str | None = None, + ) -> int: + """按用户+状态筛选计数。""" + query = self.session.query(GenerationTaskModel).filter( + GenerationTaskModel.created_by_user_id == user_id + ) + if status: + query = query.filter(GenerationTaskModel.status == status) + return query.count() + + def list_by_project_filtered( + self, + project_id: str, + *, + status: str | None = None, + limit: int | None = None, + offset: int = 0, + ) -> list[GenerationTask]: + """按项目+状态筛选任务列表。""" + query = self.session.query(GenerationTaskModel).filter( + GenerationTaskModel.project_id == project_id + ) + if status: + query = query.filter(GenerationTaskModel.status == status) + query = query.order_by(GenerationTaskModel.created_at.desc()) + if offset: + query = query.offset(offset) + if limit: + query = query.limit(limit) + return [_to_domain(m) for m in query.all()] + + def count_by_project_filtered( + self, + project_id: str, + *, + status: str | None = None, + ) -> int: + """按项目+状态筛选计数。""" + query = self.session.query(GenerationTaskModel).filter( + GenerationTaskModel.project_id == project_id + ) + if status: + query = query.filter(GenerationTaskModel.status == status) + return query.count() + def update(self, task: GenerationTask) -> GenerationTask: model = self.session.query(GenerationTaskModel).filter(GenerationTaskModel.id == task.id).first() if model is None: @@ -143,6 +221,10 @@ class SQLAlchemyGenerationTaskRepository: model.progress = task.progress model.result_count = task.result_count model.error_message = task.error_message + model.error_info = task.error_info or None + model.retry_count = task.retry_count or 0 + model.auto_retry_enabled = task.auto_retry_enabled + model.auto_retry_max = task.auto_retry_max or 0 model.started_at = task.started_at model.completed_at = task.completed_at model.source_edit_plan_id = task.source_edit_plan_id or None diff --git a/packages/adapters/sqlalchemy_impl/models.py b/packages/adapters/sqlalchemy_impl/models.py index dbfed0d2e..fb955846f 100755 --- a/packages/adapters/sqlalchemy_impl/models.py +++ b/packages/adapters/sqlalchemy_impl/models.py @@ -250,6 +250,10 @@ class GenerationTaskModel(Base): progress = Column(Float, nullable=False, default=0.0) result_count = Column(Float, nullable=False, default=0) error_message = Column(Text, nullable=False, default="") + error_info = Column(JSON, nullable=True) + retry_count = Column(Integer, nullable=False, default=0) + auto_retry_enabled = Column(Boolean, nullable=False, default=False) + auto_retry_max = Column(Integer, nullable=False, default=0) started_at = Column(DateTime, nullable=True) completed_at = Column(DateTime, nullable=True) created_by_user_id = Column(String(36), nullable=False, default="", index=True) diff --git a/packages/application/__init__.py b/packages/application/__init__.py index e0c234713..61770f3aa 100755 --- a/packages/application/__init__.py +++ b/packages/application/__init__.py @@ -28,6 +28,9 @@ from .generation_tasks import ( CreateGenerationTaskCommand, CreateGenerationTaskUseCase, GetGenerationTaskUseCase, + ListGenerationTasksResult, + ListUserTasksFilteredUseCase, + RetryGenerationTaskUseCase, ) from .ingest_jobs import SubmitIngestJobCommand, SubmitIngestJobUseCase from .jobs import ( @@ -65,6 +68,9 @@ __all__ = [ "CreateGenerationTaskCommand", "CreateGenerationTaskUseCase", "GetGenerationTaskUseCase", + "ListGenerationTasksResult", + "ListUserTasksFilteredUseCase", + "RetryGenerationTaskUseCase", "CreateJobCommand", "CreateJobUseCase", "CreateProjectCommand", diff --git a/packages/application/generation_tasks.py b/packages/application/generation_tasks.py old mode 100644 new mode 100755 index 894b44032..0e99812ee --- a/packages/application/generation_tasks.py +++ b/packages/application/generation_tasks.py @@ -21,6 +21,8 @@ class CreateGenerationTaskCommand: source_edit_plan_id: str = "" asset_select_mode: str = "" batch_id: str = "" + auto_retry_enabled: bool = False + auto_retry_max: int = 0 class CreateGenerationTaskUseCase: @@ -42,12 +44,12 @@ class CreateGenerationTaskUseCase: progress=0.0, result_count=0, error_message="", - started_at=None, - completed_at=None, created_by_user_id=command.created_by_user_id, source_edit_plan_id=command.source_edit_plan_id, asset_select_mode=command.asset_select_mode, batch_id=command.batch_id, + auto_retry_enabled=command.auto_retry_enabled, + auto_retry_max=command.auto_retry_max, ) return self.generation_task_repository.create(task) @@ -58,3 +60,63 @@ class GetGenerationTaskUseCase: def execute(self, task_id: str) -> GenerationTask | None: return self.generation_task_repository.get(task_id) + + +@dataclass(slots=True) +class ListTasksFilter: + """任务列表筛选条件。""" + status: str | None = None # pending, running, completed, failed, cancelled + + +@dataclass(slots=True) +class ListGenerationTasksResult: + """带筛选和分页的任务列表结果。""" + items: list[GenerationTask] + total: int + + +class ListUserTasksFilteredUseCase: + """按用户+筛选条件查询任务列表。""" + + def __init__(self, generation_task_repository: GenerationTaskRepository): + self.generation_task_repository = generation_task_repository + + def execute( + self, + user_id: str, + *, + status: str | None = None, + limit: int | None = None, + offset: int = 0, + ) -> ListGenerationTasksResult: + items = self.generation_task_repository.list_by_user_filtered( + user_id, + status=status, + limit=limit, + offset=offset, + ) + total = self.generation_task_repository.count_by_user_filtered( + user_id, + status=status, + ) + return ListGenerationTasksResult(items=items, total=total) + + +class RetryGenerationTaskUseCase: + """原地重试失败的任务(重置状态+递增retry_count)。 + + 与创建新任务不同:复用同一个 task_id,保留历史关联。 + """ + + def __init__(self, generation_task_repository: GenerationTaskRepository): + self.generation_task_repository = generation_task_repository + + def execute(self, task_id: str) -> GenerationTask: + task = self.generation_task_repository.get(task_id) + if task is None: + raise ValueError(f"任务不存在: {task_id}") + if not task.is_failed: + raise ValueError(f"只有失败状态的任务才能重试,当前状态: {task.status.value}") + task.mark_pending_from_failed() + self.generation_task_repository.update(task) + return task diff --git a/packages/domain/generation_task.py b/packages/domain/generation_task.py old mode 100644 new mode 100755 index 7222add07..aced09ea3 --- a/packages/domain/generation_task.py +++ b/packages/domain/generation_task.py @@ -80,6 +80,10 @@ class GenerationTask: progress: float = 0.0 result_count: int = 0 error_message: str = "" + error_info: dict = field(default_factory=dict) + retry_count: int = 0 + auto_retry_enabled: bool = False + auto_retry_max: int = 0 started_at: datetime | None = None completed_at: datetime | None = None source_edit_plan_id: str = "" @@ -105,6 +109,8 @@ class GenerationTask: source_edit_plan_id: str = "", asset_select_mode: str = "", batch_id: str = "", + auto_retry_enabled: bool = False, + auto_retry_max: int = 0, ) -> "GenerationTask": if not project_id.strip() and not template_id.strip(): raise ValueError("project_id 或 template_id 至少需要提供一个") @@ -124,6 +130,8 @@ class GenerationTask: source_edit_plan_id=source_edit_plan_id.strip(), asset_select_mode=asset_select_mode, batch_id=batch_id, + auto_retry_enabled=auto_retry_enabled, + auto_retry_max=auto_retry_max, ) # ── 状态查询 ──────────────────────────────────────────────────────────── @@ -203,13 +211,14 @@ class GenerationTask: self.result_count = result_count self.error_message = "" - def mark_failed(self, error_message: str) -> None: + def mark_failed(self, error_message: str, error_info: dict | None = None) -> None: """标记为失败(pending / running → failed)。 - 设置 error_message、completed_at。 + 设置 error_message、error_info、completed_at。 Args: error_message: 错误信息 + error_info: 结构化错误信息(error_type, stack_trace, stage, failed_at等) Raises: ValueError: 当前状态不允许转换到 failed @@ -217,6 +226,14 @@ class GenerationTask: self.transition_to(GenerationTaskStatus.FAILED) self.error_message = error_message self.completed_at = datetime.now(timezone.utc) + if error_info is not None: + self.error_info = error_info + else: + self.error_info = { + "error_type": "UnknownError", + "message": error_message, + "failed_at": datetime.now(timezone.utc).isoformat(), + } def mark_cancelled(self) -> None: """标记为已取消(pending / running → cancelled)。 @@ -269,7 +286,8 @@ class GenerationTask: def mark_pending_from_failed(self) -> None: """从失败状态重置为待处理(用于重试)。 - 清除 error_message、started_at、completed_at、progress。 + 清除 error_message、error_info、started_at、completed_at、progress, + 递增 retry_count。 Raises: ValueError: 当前状态不是 failed @@ -278,7 +296,9 @@ class GenerationTask: raise ValueError(f"只有 failed 状态的任务可以重置为 pending,当前状态: {self.status.value}") self.transition_to(GenerationTaskStatus.PENDING) self.error_message = "" + self.error_info = {} self.started_at = None self.completed_at = None self.progress = 0.0 self.result_count = 0 + self.retry_count += 1 diff --git a/packages/ports/generation_task_repository.py b/packages/ports/generation_task_repository.py index a86314aa9..5c21200e5 100755 --- a/packages/ports/generation_task_repository.py +++ b/packages/ports/generation_task_repository.py @@ -24,4 +24,36 @@ class GenerationTaskRepository(Protocol): def list_by_source_edit_plan(self, plan_id: str) -> list[GenerationTask]: ... + def list_by_user_filtered( + self, + user_id: str, + *, + status: str | None = None, + limit: int | None = None, + offset: int = 0, + ) -> list[GenerationTask]: ... + + def count_by_user_filtered( + self, + user_id: str, + *, + status: str | None = None, + ) -> int: ... + + def list_by_project_filtered( + self, + project_id: str, + *, + status: str | None = None, + limit: int | None = None, + offset: int = 0, + ) -> list[GenerationTask]: ... + + def count_by_project_filtered( + self, + project_id: str, + *, + status: str | None = None, + ) -> int: ... + def update(self, task: GenerationTask) -> GenerationTask: ... diff --git a/tests/unit/test_generation_task_status.py b/tests/unit/test_generation_task_status.py old mode 100644 new mode 100755 index dd46815e6..d2d2dd04d --- a/tests/unit/test_generation_task_status.py +++ b/tests/unit/test_generation_task_status.py @@ -453,3 +453,105 @@ class TestFullFlow: task.mark_cancelled() assert task.status == GenerationTaskStatus.CANCELLED assert task.is_terminal + + +# ── 错误信息与重试(任务中心升级) ────────────────────────────────────────── + + +class TestErrorInfo: + """测试 error_info 结构化错误信息。""" + + def test_mark_failed_default_error_info(self) -> None: + """mark_failed 不传 error_info 时自动生成默认结构。""" + task = _make_task() + task.mark_processing() + task.mark_failed("something went wrong") + assert task.is_failed + assert task.error_message == "something went wrong" + assert task.error_info["error_type"] == "UnknownError" + assert task.error_info["message"] == "something went wrong" + assert "failed_at" in task.error_info + + def test_mark_failed_with_custom_error_info(self) -> None: + """mark_failed 传自定义 error_info。""" + task = _make_task() + task.mark_processing() + info = { + "error_type": "FFmpegError", + "message": "Invalid data found", + "stack_trace": "Traceback...", + "stage": "render", + "failed_at": "2026-01-01T00:00:00+00:00", + } + task.mark_failed("Invalid data found", error_info=info) + assert task.error_info == info + + def test_error_info_cleared_on_retry(self) -> None: + """重试时 error_info 被清空。""" + task = _make_task() + task.mark_processing() + task.mark_failed("oops") + assert task.error_info # 失败时有值 + task.mark_pending_from_failed() + assert task.error_info == {} + assert task.status == GenerationTaskStatus.PENDING + + +class TestRetryCount: + """测试 retry_count 重试次数。""" + + def test_default_retry_count_is_zero(self) -> None: + """新任务 retry_count 默认 0。""" + task = _make_task() + assert task.retry_count == 0 + + def test_retry_increments_count(self) -> None: + """每次失败后重试,retry_count +1。""" + task = _make_task() + task.mark_processing() + task.mark_failed("fail 1") + task.mark_pending_from_failed() + assert task.retry_count == 1 + + task.mark_processing() + task.mark_failed("fail 2") + task.mark_pending_from_failed() + assert task.retry_count == 2 + + def test_completed_does_not_affect_retry_count(self) -> None: + """正常完成不改变 retry_count。""" + task = _make_task() + task.mark_processing() + task.mark_completed() + assert task.retry_count == 0 + + +class TestAutoRetryConfig: + """测试自动重试配置。""" + + def test_default_auto_retry_disabled(self) -> None: + """默认关闭自动重试。""" + task = _make_task() + assert task.auto_retry_enabled is False + assert task.auto_retry_max == 0 + + def test_create_with_auto_retry(self) -> None: + """create 工厂方法支持 auto_retry 参数。""" + task = GenerationTask.create( + project_id="proj-1", + asset_library_id="lib-1", + auto_retry_enabled=True, + auto_retry_max=3, + ) + assert task.auto_retry_enabled is True + assert task.auto_retry_max == 3 + + def test_auto_retry_max_default_zero(self) -> None: + """auto_retry_max 默认 0 表示不自动重试。""" + task = GenerationTask.create( + project_id="proj-1", + asset_library_id="lib-1", + auto_retry_enabled=True, + ) + assert task.auto_retry_enabled is True + assert task.auto_retry_max == 0 -- 2.54.0 From ce65e9439b7d41d6c95cd761ec6b4dfeafccc933 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Tue, 14 Jul 2026 09:17:39 +0800 Subject: [PATCH 2/6] =?UTF-8?q?chore:=20=E4=BF=AE=E5=A4=8Dblack/isort?= =?UTF-8?q?=E6=A0=BC=E5=BC=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/api/app/schemas/generation_task.py | 4 +++- apps/worker/worker_app/tasks/generation.py | 14 ++++++++++---- .../generation_task_repository.py | 16 ++++------------ packages/application/generation_tasks.py | 2 ++ 4 files changed, 19 insertions(+), 17 deletions(-) diff --git a/apps/api/app/schemas/generation_task.py b/apps/api/app/schemas/generation_task.py index 4f7a812a4..10543d0f1 100755 --- a/apps/api/app/schemas/generation_task.py +++ b/apps/api/app/schemas/generation_task.py @@ -39,7 +39,9 @@ class CreateGenerationTaskRequest(BaseModel): description="是否开启失败自动重试,默认关闭", ) auto_retry_max: int = Field( - default=0, ge=0, le=5, + default=0, + ge=0, + le=5, description="最大自动重试次数,0表示不自动重试,最大5次", ) diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index 7b59b91dd..2ea48cbf8 100755 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -1197,10 +1197,12 @@ def generate_video(self, task_id: str) -> dict: if current_retry < _task.auto_retry_max: logger.info( "[task_id=%s] 触发自动重试: 当前重试次数=%d, 最大重试次数=%d", - task_id, current_retry, _task.auto_retry_max, + task_id, + current_retry, + _task.auto_retry_max, ) # 计算退避延迟(指数退避,基础5s,最大60s) - backoff_seconds = min(5 * (2 ** current_retry), 60) + backoff_seconds = min(5 * (2**current_retry), 60) # 原地重试 _task.mark_pending_from_failed() _r.update(_task) @@ -1212,14 +1214,18 @@ def generate_video(self, task_id: str) -> dict: ) logger.info( "[task_id=%s] 自动重试已入队: 延迟=%ds, 第%d次重试", - task_id, backoff_seconds, current_retry + 1, + task_id, + backoff_seconds, + current_retry + 1, ) finally: _s.close() except Exception as retry_err: logger.warning( "[task_id=%s] 自动重试逻辑执行失败: %s", - task_id, retry_err, exc_info=True, + task_id, + retry_err, + exc_info=True, ) return { diff --git a/packages/adapters/sqlalchemy_impl/generation_task_repository.py b/packages/adapters/sqlalchemy_impl/generation_task_repository.py index 53c1fa31c..6d26a2d54 100755 --- a/packages/adapters/sqlalchemy_impl/generation_task_repository.py +++ b/packages/adapters/sqlalchemy_impl/generation_task_repository.py @@ -144,9 +144,7 @@ class SQLAlchemyGenerationTaskRepository: offset: int = 0, ) -> list[GenerationTask]: """按用户+状态筛选任务列表。""" - query = self.session.query(GenerationTaskModel).filter( - GenerationTaskModel.created_by_user_id == user_id - ) + query = self.session.query(GenerationTaskModel).filter(GenerationTaskModel.created_by_user_id == user_id) if status: query = query.filter(GenerationTaskModel.status == status) query = query.order_by(GenerationTaskModel.created_at.desc()) @@ -163,9 +161,7 @@ class SQLAlchemyGenerationTaskRepository: status: str | None = None, ) -> int: """按用户+状态筛选计数。""" - query = self.session.query(GenerationTaskModel).filter( - GenerationTaskModel.created_by_user_id == user_id - ) + query = self.session.query(GenerationTaskModel).filter(GenerationTaskModel.created_by_user_id == user_id) if status: query = query.filter(GenerationTaskModel.status == status) return query.count() @@ -179,9 +175,7 @@ class SQLAlchemyGenerationTaskRepository: offset: int = 0, ) -> list[GenerationTask]: """按项目+状态筛选任务列表。""" - query = self.session.query(GenerationTaskModel).filter( - GenerationTaskModel.project_id == project_id - ) + query = self.session.query(GenerationTaskModel).filter(GenerationTaskModel.project_id == project_id) if status: query = query.filter(GenerationTaskModel.status == status) query = query.order_by(GenerationTaskModel.created_at.desc()) @@ -198,9 +192,7 @@ class SQLAlchemyGenerationTaskRepository: status: str | None = None, ) -> int: """按项目+状态筛选计数。""" - query = self.session.query(GenerationTaskModel).filter( - GenerationTaskModel.project_id == project_id - ) + query = self.session.query(GenerationTaskModel).filter(GenerationTaskModel.project_id == project_id) if status: query = query.filter(GenerationTaskModel.status == status) return query.count() diff --git a/packages/application/generation_tasks.py b/packages/application/generation_tasks.py index 0e99812ee..f9738dfd0 100755 --- a/packages/application/generation_tasks.py +++ b/packages/application/generation_tasks.py @@ -65,12 +65,14 @@ class GetGenerationTaskUseCase: @dataclass(slots=True) class ListTasksFilter: """任务列表筛选条件。""" + status: str | None = None # pending, running, completed, failed, cancelled @dataclass(slots=True) class ListGenerationTasksResult: """带筛选和分页的任务列表结果。""" + items: list[GenerationTask] total: int -- 2.54.0 From 900fbad18e5db05bd88c090c8b4eb31c27d18d4f Mon Sep 17 00:00:00 2001 From: CI Bot Date: Tue, 14 Jul 2026 09:26:37 +0800 Subject: [PATCH 3/6] =?UTF-8?q?chore:=20=E6=A0=BC=E5=BC=8F=E5=8C=96alembic?= =?UTF-8?q?=E8=BF=81=E7=A7=BB=E6=96=87=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...rror_info_and_retry_to_generation_tasks.py | 29 ++++++++++--------- 1 file changed, 15 insertions(+), 14 deletions(-) diff --git a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py index 3005c83c6..4f8a4ef3a 100644 --- a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py +++ b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py @@ -4,13 +4,14 @@ Revision ID: 038 Revises: 037 Create Date: 2026-07-13 22:15:00.000000 """ + from alembic import op import sqlalchemy as sa from sqlalchemy.dialects.mysql import JSON as MySQLJSON # revision identifiers, used by Alembic. -revision = '038' -down_revision = '037' +revision = "038" +down_revision = "037" branch_labels = None depends_on = None @@ -18,28 +19,28 @@ depends_on = None def upgrade(): # error_info: 结构化错误信息(error_type, message, stack_trace, failed_at, stage等) op.add_column( - 'generation_tasks', - sa.Column('error_info', sa.JSON(), nullable=True), + "generation_tasks", + sa.Column("error_info", sa.JSON(), nullable=True), ) # retry_count: 重试次数 op.add_column( - 'generation_tasks', - sa.Column('retry_count', sa.Integer(), nullable=False, server_default='0'), + "generation_tasks", + sa.Column("retry_count", sa.Integer(), nullable=False, server_default="0"), ) # auto_retry_enabled: 是否开启自动重试 op.add_column( - 'generation_tasks', - sa.Column('auto_retry_enabled', sa.Boolean(), nullable=False, server_default=sa.text('false')), + "generation_tasks", + sa.Column("auto_retry_enabled", sa.Boolean(), nullable=False, server_default=sa.text("false")), ) # auto_retry_max: 最大自动重试次数 op.add_column( - 'generation_tasks', - sa.Column('auto_retry_max', sa.Integer(), nullable=False, server_default='0'), + "generation_tasks", + sa.Column("auto_retry_max", sa.Integer(), nullable=False, server_default="0"), ) def downgrade(): - op.drop_column('generation_tasks', 'auto_retry_max') - op.drop_column('generation_tasks', 'auto_retry_enabled') - op.drop_column('generation_tasks', 'retry_count') - op.drop_column('generation_tasks', 'error_info') + op.drop_column("generation_tasks", "auto_retry_max") + op.drop_column("generation_tasks", "auto_retry_enabled") + op.drop_column("generation_tasks", "retry_count") + op.drop_column("generation_tasks", "error_info") -- 2.54.0 From 5c582e0983a2a7cf6ca180546b87522d2a3f0715 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Tue, 14 Jul 2026 09:26:48 +0800 Subject: [PATCH 4/6] =?UTF-8?q?chore:=20=E4=BF=AE=E5=A4=8Dalembic=E7=9A=84?= =?UTF-8?q?isort=E6=8E=92=E5=BA=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../038_add_error_info_and_retry_to_generation_tasks.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py index 4f8a4ef3a..aba2518d1 100644 --- a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py +++ b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py @@ -5,10 +5,11 @@ Revises: 037 Create Date: 2026-07-13 22:15:00.000000 """ -from alembic import op import sqlalchemy as sa from sqlalchemy.dialects.mysql import JSON as MySQLJSON +from alembic import op + # revision identifiers, used by Alembic. revision = "038" down_revision = "037" -- 2.54.0 From 6c58e14951bcf4175854c26a1811bbad2a6cfb52 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Tue, 14 Jul 2026 09:41:37 +0800 Subject: [PATCH 5/6] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8Dalembic=E8=BF=81?= =?UTF-8?q?=E7=A7=BB=E7=89=88=E6=9C=AC=E5=8F=B7=E4=B8=8Edevelop=E4=B8=8D?= =?UTF-8?q?=E5=8C=B9=E9=85=8D=E7=9A=84=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../038_add_error_info_and_retry_to_generation_tasks.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py index aba2518d1..02ee2cef5 100644 --- a/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py +++ b/alembic/versions/038_add_error_info_and_retry_to_generation_tasks.py @@ -1,7 +1,7 @@ """add error_info and retry fields to generation_tasks -Revision ID: 038 -Revises: 037 +Revision ID: 038_error_retry +Revises: 037_generation_logs Create Date: 2026-07-13 22:15:00.000000 """ @@ -11,8 +11,8 @@ from sqlalchemy.dialects.mysql import JSON as MySQLJSON from alembic import op # revision identifiers, used by Alembic. -revision = "038" -down_revision = "037" +revision = "038_error_retry" +down_revision = "037_generation_logs" branch_labels = None depends_on = None -- 2.54.0 From 761b21fbb65bcb33fc63c0921b85e49597ec71f8 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Tue, 14 Jul 2026 09:53:05 +0800 Subject: [PATCH 6/6] =?UTF-8?q?chore:=20=E6=9B=B4=E6=96=B0schema=20metadat?= =?UTF-8?q?a=20snapshot?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/schema-metadata-snapshot.json | 32 ++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/docs/schema-metadata-snapshot.json b/docs/schema-metadata-snapshot.json index f500bc075..45a76d430 100644 --- a/docs/schema-metadata-snapshot.json +++ b/docs/schema-metadata-snapshot.json @@ -1493,6 +1493,38 @@ "type": "TEXT", "unique": false }, + { + "index": false, + "name": "error_info", + "nullable": true, + "primary_key": false, + "type": "JSON", + "unique": false + }, + { + "index": false, + "name": "retry_count", + "nullable": false, + "primary_key": false, + "type": "INTEGER", + "unique": false + }, + { + "index": false, + "name": "auto_retry_enabled", + "nullable": false, + "primary_key": false, + "type": "BOOLEAN", + "unique": false + }, + { + "index": false, + "name": "auto_retry_max", + "nullable": false, + "primary_key": false, + "type": "INTEGER", + "unique": false + }, { "index": false, "name": "started_at", -- 2.54.0