feat(worker): celery 队列隔离 + 孤儿任务消息作废 (#1714) #1722

Merged
xiaoxia merged 2 commits from feature/celery-queue-isolation-1714 into develop 2026-09-05 20:10:13 +08:00
Owner

背景(Issue #1714 补充任务)

  1. Celery 队列隔离:素材转码 ingest_asset 与视频生成共用默认队列、worker 单进程消费,实测 21 个转码积压把用户生成任务堵 40+ 分钟。
  2. 孤儿恢复消息作废缺陷:清理把 pending 任务标 failed 后未作废 Redis 队列消息,重启后消息重投 → failed → running 非法状态转换,worker 打 ERROR 后仍跑完产出半成品。

改动

队列隔离(三队列拓扑)

  • 新增 packages/shared/celery_queues.pygeneration / transcode / celery 三队列 + task_routes
    • worker.generate_videogeneration(用户视频生成,高优先级)
    • ingest_asset / classify_asset / process_duplication_check / check_duplicatetranscode(素材入库后处理链路)
    • tts/voice/batch_download/cleanup 等 → celery 默认队列
  • worker 入口 entrypoint-worker.sh双进程
    • generation worker:-Q generation -B(beat 内嵌,单实例约束),worker_prefetch_multiplier=1task_acks_late=True,并发 GENERATION_CONCURRENCY 默认 2
    • transcode worker:-Q transcode,celery,并发=总并发-2(最小 1)
    • wait -n 任一退出则 TERM 另一个;compose / deploy-staging / deploy-production / start-worker.ps1 同步
  • 双进程从进程级保证:transcode 积压 20+ 时 generation 仍有独立并发槽立即领取

消息作废(双保险)

  • 新增 packages/shared/celery_orphan_guard.py
    • revoke_and_purge():control.revoke 广播(在线 worker)+ 直接扫 Redis list 物理移除消息体(worker 下线期间 revoke 收不到);解析 celery JSON 信封,按业务 id + celery headers.id 双匹配,未命中 rpush 回写保序,解析失败保守保留
    • ensure_task_claimable():终态(failed/cancelled/completed)抛 StaleTaskDiscarded
  • 入队持久化 celery 消息 id:send_task 返回 id 写入 generation_tasks.celery_task_id / ingest_jobs.celery_task_id(新列,067 迁移,server_default 空串,写入失败仅 warning 不阻断);覆盖生成/上传/分片/重试全部入队点
  • 执行前 DB 状态守卫:generate_video / ingest_asset 加载任务后终态直接返回 {"status":"discarded"} 不进业务逻辑;mark_processing 返回 False(非法转换如 failed→running)安全中止,不再打 ERROR 后继续渲染
  • 孤儿/超时清理(_startup 启动恢复 + cleanup beat)标 failed 时同步 revoke + 清队列
  • pending 超时阈值 15→45 分钟,与 running 孤儿 20min 区分,避免正常排队误杀(队列隔离后 generation 排队极短,45min 仅兜底 worker 停消费)

测试

  • 新增 22 个单测:
    • 路由表 / apply_queue_settings / Redis 真实消息清理(biz id + celery id 双匹配)/ revoke 调用 / 终态守卫各状态
    • 端到端:清理标 failed 后队列消息被移除、正常消息保留不重投
    • ingest failed/completed 丢弃不 commit、generation failed 丢弃不渲染、claim 失败安全中止
    • 入队后 celery_task_id 持久化
  • 全量单测 14301 passed(5 failed 为预存在 h2.config TTS 环境问题,与本次无关)
  • 067 迁移隔离 DDL 验证:upgrade 加列 + server_default、downgrade 删列均正确;迁移链 67 版完整
  • black/isort/ruff/mypy 门禁通过

staging 验证方式

部署后:docker exec xiaoxia-worker-staging pgrep -af celery 应见 generation@/transcode@ 两进程;celery inspect active_queues 确认消费队列分离;向 transcode 灌入 20+ 任务后提交生成任务,1 分钟内开始渲染。

配套

  • 转码任务重试上限/失败兜底:灵应(并行)
  • WECHAT 环境变量等不涉及本 PR
## 背景(Issue #1714 补充任务) 1. **Celery 队列隔离**:素材转码 `ingest_asset` 与视频生成共用默认队列、worker 单进程消费,实测 21 个转码积压把用户生成任务堵 40+ 分钟。 2. **孤儿恢复消息作废缺陷**:清理把 pending 任务标 failed 后未作废 Redis 队列消息,重启后消息重投 → `failed → running` 非法状态转换,worker 打 ERROR 后仍跑完产出半成品。 ## 改动 ### 队列隔离(三队列拓扑) - 新增 `packages/shared/celery_queues.py`:`generation` / `transcode` / `celery` 三队列 + `task_routes` - `worker.generate_video` → **generation**(用户视频生成,高优先级) - `ingest_asset` / `classify_asset` / `process_duplication_check` / `check_duplicate` → **transcode**(素材入库后处理链路) - tts/voice/batch_download/cleanup 等 → celery 默认队列 - worker 入口 `entrypoint-worker.sh` 改**双进程**: - generation worker:`-Q generation -B`(beat 内嵌,单实例约束),`worker_prefetch_multiplier=1`、`task_acks_late=True`,并发 `GENERATION_CONCURRENCY` 默认 2 - transcode worker:`-Q transcode,celery`,并发=总并发-2(最小 1) - `wait -n` 任一退出则 TERM 另一个;compose / deploy-staging / deploy-production / start-worker.ps1 同步 - 双进程从进程级保证:transcode 积压 20+ 时 generation 仍有独立并发槽立即领取 ### 消息作废(双保险) - 新增 `packages/shared/celery_orphan_guard.py`: - `revoke_and_purge()`:control.revoke 广播(在线 worker)+ 直接扫 Redis list 物理移除消息体(worker 下线期间 revoke 收不到);解析 celery JSON 信封,按**业务 id + celery headers.id 双匹配**,未命中 rpush 回写保序,解析失败保守保留 - `ensure_task_claimable()`:终态(failed/cancelled/completed)抛 StaleTaskDiscarded - **入队持久化 celery 消息 id**:send_task 返回 id 写入 `generation_tasks.celery_task_id` / `ingest_jobs.celery_task_id`(新列,067 迁移,server_default 空串,写入失败仅 warning 不阻断);覆盖生成/上传/分片/重试全部入队点 - **执行前 DB 状态守卫**:generate_video / ingest_asset 加载任务后终态直接返回 `{"status":"discarded"}` 不进业务逻辑;`mark_processing` 返回 False(非法转换如 failed→running)安全中止,不再打 ERROR 后继续渲染 - 孤儿/超时清理(_startup 启动恢复 + cleanup beat)标 failed 时同步 revoke + 清队列 - pending 超时阈值 **15→45 分钟**,与 running 孤儿 20min 区分,避免正常排队误杀(队列隔离后 generation 排队极短,45min 仅兜底 worker 停消费) ## 测试 - 新增 22 个单测: - 路由表 / apply_queue_settings / Redis 真实消息清理(biz id + celery id 双匹配)/ revoke 调用 / 终态守卫各状态 - 端到端:清理标 failed 后队列消息被移除、正常消息保留不重投 - ingest failed/completed 丢弃不 commit、generation failed 丢弃不渲染、claim 失败安全中止 - 入队后 celery_task_id 持久化 - 全量单测 **14301 passed**(5 failed 为预存在 h2.config TTS 环境问题,与本次无关) - 067 迁移隔离 DDL 验证:upgrade 加列 + server_default、downgrade 删列均正确;迁移链 67 版完整 - black/isort/ruff/mypy 门禁通过 ## staging 验证方式 部署后:`docker exec xiaoxia-worker-staging pgrep -af celery` 应见 generation@/transcode@ 两进程;celery inspect active_queues 确认消费队列分离;向 transcode 灌入 20+ 任务后提交生成任务,1 分钟内开始渲染。 ## 配套 - 转码任务重试上限/失败兜底:灵应(并行) - WECHAT 环境变量等不涉及本 PR

🚀 预览环境已部署

项目 详情
PR号 #1722
预览链接 https://pr-1722.preview.xiaoxiajianji.com
API环境 staging

💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。

🔄 每次提交新代码后预览环境会自动更新。

🗑️ PR 关闭或合并后,预览环境会自动清理。

🚀 **预览环境已部署** | 项目 | 详情 | |------|------| | PR号 | #1722 | | 预览链接 | [https://pr-1722.preview.xiaoxiajianji.com](https://pr-1722.preview.xiaoxiajianji.com) | | API环境 | staging | > 💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。 > > 🔄 每次提交新代码后预览环境会自动更新。 > > 🗑️ PR 关闭或合并后,预览环境会自动清理。
xiaoxia added 1 commit 2026-09-05 19:11:33 +08:00
feat(worker): celery 队列隔离 + 孤儿任务消息作废 (#1714)
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
df99305dd6
问题:素材转码与视频生成共用 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 通过。
xiaoxia force-pushed feature/celery-queue-isolation-1714 from f5bdf17da1 to df99305dd6 2026-09-05 19:11:33 +08:00 Compare
Collaborator

【阻塞级判定】

  • 是否存在阻塞级问题:是
  • 阻塞级问题数量:3 个

📊 审查概览

  • 整体评价:需修改
  • 建议级问题数量:2 个

🔴 阻塞级问题(必须修复)

  1. [apps/api/app/api/routes/chunked_upload.py:384] 任务ID持久化失败导致孤儿任务无法清理

    • 问题类型:数据损坏风险/逻辑bug
    • 问题描述:在发送Celery任务后,代码尝试将celery_task_id持久化到数据库,但使用了裸except Exception: pass。如果数据库更新失败,celery_task_id将无法记录。当后续触发超时清理逻辑时,_revoke_and_purge_stale_messages将无法获取到有效的celery_task_id,导致无法撤销队列中的消息。这会导致DB中任务状态已标记为失败,但队列中仍存在旧消息,Worker重启后可能重新消费并执行,造成非法状态转换和资源浪费。
    • 修改建议:不应在此处静默失败。建议至少记录日志,或者考虑在任务入队前预先分配ID并持久化,或者在持久化失败时主动调用celery_app.control.revoke撤销已发送的任务。
  2. [apps/api/app/api/routes/ingest_jobs.py:46] 任务ID持久化失败导致孤儿任务无法清理

    • 问题类型:数据损坏风险/逻辑bug
    • 问题描述:同上,ingest_jobs.py中使用了相同的裸except Exception: pass模式处理celery_task_id的持久化。这会导致素材转码任务在DB更新失败时丢失与Celery消息的关联,后续清理流程无法正确撤销消息,可能引发重复处理。
    • 修改建议:同上,增加日志记录或采取回滚机制(撤销已发送的任务),确保DB状态与消息队列状态的一致性。
  3. [apps/api/app/api/routes/task_center.py:378] 任务ID持久化失败导致孤儿任务无法清理

    • 问题类型:数据损坏风险/逻辑bug
    • 问题描述:同上,task_center.py中的重试逻辑也存在相同问题。静默忽略持久化错误会导致重试生成的任务与Celery消息断链,破坏了#1714 issue修复的核心目标(通过ID关联进行消息清理)。
    • 修改建议:同上,必须处理持久化失败的情况,不能静默忽略。

💡 改进建议(不阻塞合并)

  1. [apps/worker/worker_app/tasks/_startup.py:45] 降级逻辑可能掩盖真实错误

    • 具体内容:cleanup_stale_running_with_session_ids中,当repo不支持_with_ids方法时,降级返回[("", "")]。虽然这兼容了旧代码,但如果cleanup_stale_running本身抛出异常(如数据库连接问题),这里会直接崩溃,且没有日志提示“使用了降级模式”。建议在降级分支增加日志记录,方便排查问题。
  2. [packages/adapters/sqlalchemy_impl/generation_task_repository.py:329] 批量更新改为循环更新影响性能

    • 具体内容:cleanup_stale_pending_with_ids方法中,原代码使用了.update(..., synchronize_session=False)进行批量更新,效率较高。修改后的代码先.all()查出所有对象,再循环修改属性并提交。如果待清理的pending任务数量较多(例如积压场景),这会一次性加载大量数据到内存,且产生大量的UPDATE语句。建议在数据量可控的情况下保留,或者考虑分批处理以平衡内存和性能。

良好实践

  • 在Worker任务入口(generation.pyingest.py)增加了执行前的状态守卫检查,有效防止了僵尸消息的非法执行,这是解决并发/分布式环境下状态一致性的关键措施。
  • 数据库迁移脚本中为celery_task_id设置了server_default="",保证了非空约束下的平滑升级。
  • entrypoint-worker.sh脚本中通过trapwait -n实现了双Worker进程的协同退出管理,避免了孤儿进程。

格式检查通过 | 逻辑审查需修改 | ⚠️ 建议关注性能


🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 11:14:26 | 模型:

### 【阻塞级判定】 - 是否存在阻塞级问题:是 - 阻塞级问题数量:3 个 ### 📊 审查概览 - 整体评价:需修改 - 建议级问题数量:2 个 ### 🔴 阻塞级问题(必须修复) 1. **[apps/api/app/api/routes/chunked_upload.py:384] 任务ID持久化失败导致孤儿任务无法清理** - 问题类型:数据损坏风险/逻辑bug - 问题描述:在发送Celery任务后,代码尝试将`celery_task_id`持久化到数据库,但使用了裸`except Exception: pass`。如果数据库更新失败,`celery_task_id`将无法记录。当后续触发超时清理逻辑时,`_revoke_and_purge_stale_messages`将无法获取到有效的`celery_task_id`,导致无法撤销队列中的消息。这会导致DB中任务状态已标记为失败,但队列中仍存在旧消息,Worker重启后可能重新消费并执行,造成非法状态转换和资源浪费。 - 修改建议:不应在此处静默失败。建议至少记录日志,或者考虑在任务入队前预先分配ID并持久化,或者在持久化失败时主动调用`celery_app.control.revoke`撤销已发送的任务。 2. **[apps/api/app/api/routes/ingest_jobs.py:46] 任务ID持久化失败导致孤儿任务无法清理** - 问题类型:数据损坏风险/逻辑bug - 问题描述:同上,`ingest_jobs.py`中使用了相同的裸`except Exception: pass`模式处理`celery_task_id`的持久化。这会导致素材转码任务在DB更新失败时丢失与Celery消息的关联,后续清理流程无法正确撤销消息,可能引发重复处理。 - 修改建议:同上,增加日志记录或采取回滚机制(撤销已发送的任务),确保DB状态与消息队列状态的一致性。 3. **[apps/api/app/api/routes/task_center.py:378] 任务ID持久化失败导致孤儿任务无法清理** - 问题类型:数据损坏风险/逻辑bug - 问题描述:同上,`task_center.py`中的重试逻辑也存在相同问题。静默忽略持久化错误会导致重试生成的任务与Celery消息断链,破坏了#1714 issue修复的核心目标(通过ID关联进行消息清理)。 - 修改建议:同上,必须处理持久化失败的情况,不能静默忽略。 ### 💡 改进建议(不阻塞合并) 1. **[apps/worker/worker_app/tasks/_startup.py:45] 降级逻辑可能掩盖真实错误** - 具体内容:`cleanup_stale_running_with_session_ids`中,当repo不支持`_with_ids`方法时,降级返回`[("", "")]`。虽然这兼容了旧代码,但如果`cleanup_stale_running`本身抛出异常(如数据库连接问题),这里会直接崩溃,且没有日志提示“使用了降级模式”。建议在降级分支增加日志记录,方便排查问题。 2. **[packages/adapters/sqlalchemy_impl/generation_task_repository.py:329] 批量更新改为循环更新影响性能** - 具体内容:`cleanup_stale_pending_with_ids`方法中,原代码使用了`.update(..., synchronize_session=False)`进行批量更新,效率较高。修改后的代码先`.all()`查出所有对象,再循环修改属性并提交。如果待清理的pending任务数量较多(例如积压场景),这会一次性加载大量数据到内存,且产生大量的UPDATE语句。建议在数据量可控的情况下保留,或者考虑分批处理以平衡内存和性能。 ### ✅ 良好实践 - 在Worker任务入口(`generation.py`和`ingest.py`)增加了执行前的状态守卫检查,有效防止了僵尸消息的非法执行,这是解决并发/分布式环境下状态一致性的关键措施。 - 数据库迁移脚本中为`celery_task_id`设置了`server_default=""`,保证了非空约束下的平滑升级。 - `entrypoint-worker.sh`脚本中通过`trap`和`wait -n`实现了双Worker进程的协同退出管理,避免了孤儿进程。 --- ✅ 格式检查通过 | ❌ 逻辑审查需修改 | ⚠️ 建议关注性能 --- <sub>🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 11:14:26 | 模型: </sub> <!-- AI_CODE_REVIEW_AUTO_COMMENT -->
xiaoxia added 1 commit 2026-09-05 19:52:51 +08:00
test(#1714): 补 mock redis 覆盖率测试,diff coverage 98%
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 2s
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 30s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m47s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m47s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m12s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 2m40s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m11s
AI Code Review / AI Code Review (pull_request) Successful in 6m37s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 7m4s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 10m41s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 14m45s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 9s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 2m36s
6a1ec20e68
CI runner 无本地 redis,purge/revoke 与路由持久化的真实 redis 测试
全部 skip 导致这些改动零覆盖、diff coverage 45% 未达 60% 门禁:

- 新增 test_orphan_guard_purge_mocked_1714.py(31 测试,全 mock):
  _extract_business_ids 各形态(三元组/裸 args/dict args/坏 JSON/坏
  base64/空 args)、_purge_one_queue(biz/celery id 双匹配命中、未命中
  rpush 保序、全删不 rpush、解析失败保守保留、lrange/重写异常)、
  purge_stale_messages_from_queues(空 ids 早退、redis 未安装、连接失败、
  close 异常)、revoke_and_purge(逐条 revoke、异常不阻断、空 id 跳过)
- 新增 test_persist_celery_id_routes_1714.py(9 测试):ingest_jobs /
  task_center 重试 / upload helper 的 celery_task_id 持久化正常与异常
  吞掉分支、task_enqueue 持久化失败仍入队成功、celery_app 队列配置
  异常不阻断启动(用独立模块对象加载,不 reload 污染 task_enqueue)、
  仓储 update 落库 celery_task_id
- chunked_upload 完成回调去重:改为复用 upload._persist_celery_task_id
xiaoxia merged commit e7d6e69396 into develop 2026-09-05 20:10:13 +08:00

🗑️ 预览环境已清理

PR #1722 已关闭或合并,对应的预览环境已被清理。

如有需要,可以重新打开 PR 来重新生成预览环境。

🗑️ **预览环境已清理** PR #1722 已关闭或合并,对应的预览环境已被清理。 > 如有需要,可以重新打开 PR 来重新生成预览环境。
Sign in to join this conversation.