Files
xiaoxia-saas/tests/unit/test_celery_queue_isolation_1714.py
saas-backend-bot 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
feat(worker): celery 队列隔离 + 孤儿任务消息作废 (#1714)
问题:素材转码与视频生成共用 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 通过。
2026-09-05 19:07:47 +08:00

177 lines
6.1 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""#1714 队列隔离 + 作废消息清除 单元测试。
覆盖:
1. task_routesgenerate_video → generationingest_asset/classify/duplication → transcode
2. purge_stale_messages_from_queuesRedis 队列中作废任务消息被物理移除,未命中保留
3. revoke_and_purgerevoke 广播 + 队列清理同时生效
4. ensure_task_claimable:终态任务抛 StaleTaskDiscardedpending 放行
"""
from __future__ import annotations
from unittest.mock import MagicMock
import pytest
from celery import Celery
from packages.shared.celery_orphan_guard import (
StaleTaskDiscarded,
_extract_business_ids,
ensure_task_claimable,
purge_stale_messages_from_queues,
revoke_and_purge,
)
from packages.shared.celery_queues import (
QUEUE_GENERATION,
QUEUE_TRANSCODE,
apply_queue_settings,
task_routes,
)
BROKER_URL = "redis://localhost:6379/15"
TEST_QUEUES = ("_test_gen_q", "_test_transcode_q")
# ── 1. 路由表 ──────────────────────────────────────────────────────────
def test_routes_send_generation_to_generation_queue():
assert task_routes["worker.generate_video"]["queue"] == QUEUE_GENERATION
def test_routes_send_ingest_to_transcode_queue():
assert task_routes["worker.ingest_asset"]["queue"] == QUEUE_TRANSCODE
assert task_routes["worker.classify_asset"]["queue"] == QUEUE_TRANSCODE
assert task_routes["worker.process_duplication_check"]["queue"] == QUEUE_TRANSCODE
assert task_routes["worker.check_duplicate"]["queue"] == QUEUE_TRANSCODE
def test_apply_queue_settings_configures_celery_app():
app = Celery("test-routes")
apply_queue_settings(app)
queue_names = {q.name for q in app.conf.task_queues}
assert queue_names == {"generation", "transcode", "celery"}
assert app.conf.task_default_queue == "celery"
# ── Redis 队列消息清理(需要本地 redis;不可用时 skip) ─────────────────
def _redis_available() -> bool:
try:
import redis
return bool(redis.Redis.from_url(BROKER_URL).ping())
except Exception:
return False
@pytest.fixture()
def redis_client():
import redis
client = redis.Redis.from_url(BROKER_URL)
for q in TEST_QUEUES:
client.delete(q)
yield client
for q in TEST_QUEUES:
client.delete(q)
def _publish(app: Celery, queue: str, celery_id: str, business_id: str) -> None:
from kombu import Queue
from kombu.pools import producers
with app.connection_for_write() as conn:
with producers[conn].acquire(block=True) as prod:
prod.publish(
(business_id,),
exchange="",
routing_key=queue,
serializer="json",
headers={"id": celery_id, "task": "worker.generate_video"},
retry=False,
delivery_mode=1,
declare=[Queue(queue, routing_key=queue, durable=False)],
)
@pytest.mark.skipif(not _redis_available(), reason="本地 redis 不可用")
def test_purge_removes_stale_business_message_and_keeps_others(redis_client):
app = Celery("test-purge")
app.conf.broker_url = BROKER_URL
_publish(app, TEST_QUEUES[0], "celery-1", "task-KEEP-A")
_publish(app, TEST_QUEUES[0], "celery-2", "task-STALE-B")
_publish(app, TEST_QUEUES[0], "celery-3", "task-KEEP-C")
_publish(app, TEST_QUEUES[1], "celery-4", "task-STALE-B") # 同一业务任务在转码队列?不应出现但验证全队列扫描
removed = purge_stale_messages_from_queues(BROKER_URL, TEST_QUEUES, business_task_ids={"task-STALE-B"})
assert removed == 2
remaining = []
for raw in redis_client.lrange(TEST_QUEUES[0], 0, -1):
_celery_id, biz_id = _extract_business_ids(raw)
remaining.append(biz_id)
assert set(remaining) == {"task-KEEP-A", "task-KEEP-C"}
assert redis_client.llen(TEST_QUEUES[1]) == 0
@pytest.mark.skipif(not _redis_available(), reason="本地 redis 不可用")
def test_purge_matches_by_celery_message_id(redis_client):
app = Celery("test-purge-msg-id")
app.conf.broker_url = BROKER_URL
_publish(app, TEST_QUEUES[0], "celery-stale-id", "task-X")
_publish(app, TEST_QUEUES[0], "celery-good-id", "task-Y")
removed = purge_stale_messages_from_queues(BROKER_URL, TEST_QUEUES, celery_task_ids={"celery-stale-id"})
assert removed == 1
assert redis_client.llen(TEST_QUEUES[0]) == 1
@pytest.mark.skipif(not _redis_available(), reason="本地 redis 不可用")
def test_revoke_and_purge_calls_control_revoke(redis_client):
app = Celery("test-revoke")
app.conf.broker_url = BROKER_URL
app.control = MagicMock()
_publish(app, TEST_QUEUES[0], "celery-revoke-1", "task-R")
removed = revoke_and_purge(
app,
BROKER_URL,
business_task_ids={"task-R"},
celery_task_ids={"celery-revoke-1"},
queue_names=TEST_QUEUES,
)
assert removed == 1
app.control.revoke.assert_called_once_with("celery-revoke-1")
def test_purge_empty_ids_is_noop():
assert purge_stale_messages_from_queues(BROKER_URL, TEST_QUEUES) == 0
# ── 2. 执行前状态守卫 ──────────────────────────────────────────────────
def test_guard_allows_pending():
status = ensure_task_claimable("t1", lambda _id: "pending", task_label="generation")
assert status == "pending"
def test_guard_rejects_failed():
with pytest.raises(StaleTaskDiscarded) as exc:
ensure_task_claimable("t2", lambda _id: "failed", task_label="generation")
assert exc.value.task_id == "t2"
assert exc.value.status == "failed"
def test_guard_rejects_cancelled_and_completed():
with pytest.raises(StaleTaskDiscarded):
ensure_task_claimable("t3", lambda _id: "cancelled")
with pytest.raises(StaleTaskDiscarded):
ensure_task_claimable("t4", lambda _id: "completed")
def test_guard_missing_task_returns_empty():
assert ensure_task_claimable("t5", lambda _id: None) == ""