df99305dd6
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 29s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 29s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 49s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m39s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m44s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m19s
AI Code Review / AI Code Review (pull_request) Failing after 2m52s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m54s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m11s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 6m23s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 3m39s
问题:素材转码与视频生成共用 celery 默认队列、worker 单进程消费, 20+ 转码积压会把用户生成任务堵 40 分钟以上;孤儿清理把任务标 failed 后 Redis 队列消息未作废,消息被重投导致 failed→running 非法转换, worker 打印 ERROR 后继续产出半成品。 队列隔离: - 新增 packages/shared/celery_queues.py:generation/transcode/celery 三队列与 task_routes(generate_video→generation;ingest_asset/ classify_asset/duplication→transcode),apply_queue_settings() - worker 入口改双进程:generation worker 独占队列并内嵌 beat (prefetch=1, GENERATION_CONCURRENCY 默认 2),transcode worker 消费 transcode,celery(并发=总-2,最小 1),任一退出则整体终止 - compose/部署脚本/ps1 同步新增 GENERATION_CONCURRENCY 与健康检查 消息作废: - 新增 packages/shared/celery_orphan_guard.py:终态守卫 ensure_task_claimable、Redis 队列消息物理清理(JSON 信封解析, 按业务 id + celery headers.id 双匹配,未命中 rpush 保序)、 revoke_and_purge(control.revoke + 物理清队列双保险) - 入队点(生成/上传/分片/重试)send_task 后持久化 celery_task_id 到 generation_tasks/ingest_jobs(新列,067 迁移,失败仅 warning) - generate_video/ingest_asset 执行前校验 DB 状态:终态直接 discarded 不进业务逻辑;mark_processing 返回 False(非法转换)安全中止 - 孤儿/超时清理标 failed 时同时 revoke + 清队列消息 - pending 超时阈值 15→45 分钟,与 running 孤儿(20min)区分 测试:新增 22 个单测(路由表/真实 Redis 消息清理/终态守卫/ 非法转换中止/标 failed 后消息不重投/入队持久化),全量 14301 passed;067 迁移隔离 DDL 验证 upgrade/downgrade 通过。
70 lines
2.7 KiB
Python
70 lines
2.7 KiB
Python
from sqlalchemy.orm import Session
|
|
|
|
from packages.adapters.sqlalchemy_impl.models import IngestJobModel
|
|
from packages.domain import IngestJob, IngestJobStatus
|
|
|
|
|
|
class SQLAlchemyIngestJobRepository:
|
|
def __init__(self, session: Session):
|
|
self.session = session
|
|
|
|
def create(self, job: IngestJob) -> IngestJob:
|
|
model = IngestJobModel(
|
|
id=job.id,
|
|
project_id=job.project_id,
|
|
library_id=job.library_id,
|
|
storage_key=job.storage_key,
|
|
status=job.status.value,
|
|
error_message=job.error_message,
|
|
result_asset_id=job.result_asset_id,
|
|
file_hash=job.file_hash,
|
|
asset_id=job.asset_id or "",
|
|
celery_task_id=getattr(job, "celery_task_id", "") or "",
|
|
created_at=job.created_at,
|
|
updated_at=job.updated_at,
|
|
)
|
|
self.session.add(model)
|
|
self.session.commit()
|
|
return job
|
|
|
|
def get(self, job_id: str) -> IngestJob | None:
|
|
model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job_id).first()
|
|
if model is None:
|
|
return None
|
|
return IngestJob(
|
|
id=model.id,
|
|
project_id=model.project_id,
|
|
library_id=model.library_id,
|
|
storage_key=model.storage_key,
|
|
status=IngestJobStatus(model.status),
|
|
error_message=model.error_message,
|
|
result_asset_id=model.result_asset_id,
|
|
file_hash=model.file_hash or "",
|
|
asset_id=getattr(model, "asset_id", "") or "",
|
|
celery_task_id=getattr(model, "celery_task_id", "") or "",
|
|
created_at=model.created_at,
|
|
updated_at=model.updated_at,
|
|
)
|
|
|
|
def list_by_project(self, project_id: str) -> list[IngestJob]:
|
|
models = self.session.query(IngestJobModel).filter(IngestJobModel.project_id == project_id).all()
|
|
return [self.get(model.id) for model in models if self.get(model.id) is not None]
|
|
|
|
def update(self, job: IngestJob) -> IngestJob:
|
|
model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job.id).first()
|
|
if model is None:
|
|
raise ValueError(f"IngestJob {job.id} not found")
|
|
model.status = job.status.value
|
|
model.error_message = job.error_message
|
|
model.result_asset_id = job.result_asset_id
|
|
model.file_hash = job.file_hash
|
|
model.storage_key = job.storage_key
|
|
if job.asset_id:
|
|
model.asset_id = job.asset_id
|
|
celery_tid = getattr(job, "celery_task_id", "")
|
|
if celery_tid:
|
|
model.celery_task_id = celery_tid
|
|
model.updated_at = job.updated_at
|
|
self.session.commit()
|
|
return job
|