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 通过。
177 lines
6.1 KiB
Python
177 lines
6.1 KiB
Python
"""#1714 队列隔离 + 作废消息清除 单元测试。
|
||
|
||
覆盖:
|
||
1. task_routes:generate_video → generation,ingest_asset/classify/duplication → transcode
|
||
2. purge_stale_messages_from_queues:Redis 队列中作废任务消息被物理移除,未命中保留
|
||
3. revoke_and_purge:revoke 广播 + 队列清理同时生效
|
||
4. ensure_task_claimable:终态任务抛 StaleTaskDiscarded,pending 放行
|
||
"""
|
||
|
||
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) == ""
|