Files
xiaoxia-saas/tests/unit/test_stale_task_revoke_1714.py
T
xiaoxia 53fb25efcf
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Check push changed paths (push) Successful in 19s
CI/CD Pipeline / Build Staging API Image (push) Successful in 41s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 48s
CI/CD Pipeline / Integration Tests (push) Successful in 3m10s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 3m17s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m30s
CI/CD Pipeline / Validate - Style (push) Successful in 4m17s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 59s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m33s
CI/CD Pipeline / Validate - Security (push) Successful in 7m12s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m38s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m13s
CI/CD Pipeline / Unit Tests (push) Successful in 10m11s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Staging E2E Tests (push) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (push) Failing after 26h14m3s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 26h24m21s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 26h19m47s
CI/CD Pipeline / PR Build Web Image (push) Failing after 26h23m44s
CI/CD Pipeline / PR Build API Image (push) Failing after 26h23m44s
CI/CD Pipeline / Deploy Production (push) Failing after 26h13m23s
CI/CD Pipeline / Build Production Web Image (push) Failing after 26h13m26s
CI/CD Pipeline / CI Gate (push) Failing after 26h13m25s
CI/CD Pipeline / Build Production API Image (push) Failing after 26h13m26s
CI/CD Pipeline / Canary Release to Production (push) Failing after 26h13m23s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Failing after 26h19m46s
CI/CD Pipeline / Frontend Lint (push) Failing after 26h23m37s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 26h23m45s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 26h19m46s
fix(#1834): 批量修复 UP 系列静态分析警告(UP007/UP006/UP017/UP035) (#1928)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-15 12:59:17 +08:00

205 lines
7.4 KiB
Python
Raw 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.
"""Issue #1714:孤儿/超时清理标记 failed 时必须撤销并清除 Redis 队列消息。
覆盖:
- cleanup_stale_pending_with_session_ids:超时 pending 标记 failed 并返回
(task_id, celery_task_id),worker 清理流程据此 revoke + purge 队列消息
- 队列中对应业务任务的 celery 消息被物理移除(作废消息不会重投执行)
- 旧仓储(无 _with_ids 方法)降级为计数模式,不抛异常
- cleanup_stale_running_with_ids 同样返回 id 列表
"""
from __future__ import annotations
import sys
from datetime import UTC, datetime, timedelta, timezone
from pathlib import Path
from unittest.mock import MagicMock
import pytest
from sqlalchemy import create_engine, text
from sqlalchemy.orm import sessionmaker
sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "apps" / "api"))
from packages.adapters.sqlalchemy_impl.generation_task_repository import ( # noqa: E402
SQLAlchemyGenerationTaskRepository,
)
from packages.adapters.sqlalchemy_impl.models import Base # noqa: E402
from packages.domain import GenerationTask, GenerationTaskStatus # noqa: E402
BROKER_URL = "redis://localhost:6379/15"
TEST_QUEUE = "_test_revoke_q"
def _repository():
engine = create_engine("sqlite:///:memory:")
Base.metadata.create_all(engine)
session = sessionmaker(bind=engine)()
return SQLAlchemyGenerationTaskRepository(session), session, engine
def _make_task(**kwargs) -> GenerationTask:
defaults = dict(project_id="proj-1", asset_library_id="lib-1", created_by_user_id="user-1")
defaults.update(kwargs)
return GenerationTask.create(**defaults)
def _redis_available() -> bool:
try:
import redis
return bool(redis.Redis.from_url(BROKER_URL).ping())
except Exception:
return False
# ── 仓储层:返回 ids ────────────────────────────────────────────────────
def test_cleanup_stale_pending_returns_ids_with_celery_task_id():
repo, _, engine = _repository()
task = _make_task()
task.celery_task_id = "celery-msg-id-001"
repo.create(task)
# created_at 改到 60 分钟前
with engine.connect() as conn:
conn.execute(
text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"),
{"ts": datetime.now(UTC) - timedelta(minutes=60), "id": task.id},
)
conn.commit()
items = repo.cleanup_stale_pending_with_ids(timeout_minutes=45)
assert len(items) == 1
biz_id, celery_id = items[0]
assert biz_id == task.id
assert celery_id == "celery-msg-id-001"
saved = repo.get(task.id)
assert saved.status == GenerationTaskStatus.FAILED
def test_cleanup_stale_running_returns_ids():
repo, _, engine = _repository()
task = _make_task()
repo.create(task)
task.mark_processing()
task.celery_task_id = "celery-msg-id-002"
repo.update(task)
with engine.connect() as conn:
conn.execute(
text("UPDATE generation_tasks SET updated_at = :ts WHERE id = :id"),
{"ts": datetime.now(UTC) - timedelta(minutes=60), "id": task.id},
)
conn.commit()
items = repo.cleanup_stale_running_with_ids(timeout_minutes=20)
assert len(items) == 1
assert items[0][0] == task.id
assert items[0][1] == "celery-msg-id-002"
assert repo.get(task.id).status == GenerationTaskStatus.FAILED
def test_legacy_repo_without_with_ids_falls_back_to_count():
"""旧仓储只有 cleanup_stale_pending(返回 int)时降级可用,不抛异常。"""
# worker 模块加载(标准 mock 模式)
saved = set(sys.modules.keys())
mock_db = MagicMock()
mock_db.SessionLocal = MagicMock()
sys.modules["worker_app.db"] = mock_db
sys.modules["worker_app.core.config"] = MagicMock()
mock_celery = MagicMock()
mock_celery.celery_app.task = MagicMock(
side_effect=(lambda *a, **k: (a[0] if a and callable(a[0]) else (lambda f: f)))
)
sys.modules["worker_app.celery_app"] = mock_celery
worker_path = str(Path(__file__).resolve().parents[2] / "apps" / "worker")
if worker_path not in sys.path:
sys.path.insert(0, worker_path)
from worker_app.tasks import _startup # noqa: E402
class LegacyRepo:
def cleanup_stale_pending(self, timeout_minutes): # noqa: ARG002
return 3
def cleanup_stale_running(self, timeout_minutes): # noqa: ARG002
return 2
items_p = _startup.cleanup_stale_pending_with_session_ids(LegacyRepo(), 45)
items_r = _startup.cleanup_stale_running_with_session_ids(LegacyRepo(), 20)
assert len(items_p) == 3
assert len(items_r) == 2
for key in list(sys.modules.keys()):
if key not in saved and not key.startswith("video_processing"):
del sys.modules[key]
# ── 端到端:清理 → 队列消息被移除(作废消息不重投) ────────────────────
@pytest.mark.skipif(not _redis_available(), reason="本地 redis 不可用")
def test_stale_pending_cleanup_purges_redis_message():
"""任务标 failed 后,其在 Redis 队列里的 celery 消息被清除,不会被重投。"""
import redis
from celery import Celery
from kombu import Queue
from kombu.pools import producers
from packages.shared.celery_orphan_guard import purge_stale_messages_from_queues
repo, _, engine = _repository()
task = _make_task()
task.celery_task_id = "celery-stale-xyz"
repo.create(task)
with engine.connect() as conn:
conn.execute(
text("UPDATE generation_tasks SET created_at = :ts WHERE id = :id"),
{"ts": datetime.now(UTC) - timedelta(minutes=60), "id": task.id},
)
conn.commit()
# 模拟该任务的 celery 消息仍在 generation 队列里(worker 下线期间未消费)
client = redis.Redis.from_url(BROKER_URL)
client.delete(TEST_QUEUE)
app = Celery("test-e2e-revoke")
app.conf.broker_url = BROKER_URL
with app.connection_for_write() as conn:
with producers[conn].acquire(block=True) as prod:
# 作废任务消息
prod.publish(
(task.id,),
exchange="",
routing_key=TEST_QUEUE,
serializer="json",
headers={"id": "celery-stale-xyz", "task": "worker.generate_video"},
retry=False,
delivery_mode=1,
declare=[Queue(TEST_QUEUE, routing_key=TEST_QUEUE, durable=False)],
)
# 另一条正常任务消息(必须保留)
prod.publish(
("other-task-id",),
exchange="",
routing_key=TEST_QUEUE,
serializer="json",
headers={"id": "celery-keep", "task": "worker.generate_video"},
retry=False,
delivery_mode=1,
)
assert client.llen(TEST_QUEUE) == 2
# 执行清理(与 worker beat 相同流程:标 failed → 拿 ids → purge)
items = repo.cleanup_stale_pending_with_ids(timeout_minutes=45)
biz_ids = [bid for bid, _ in items]
celery_ids = [cid for _, cid in items if cid]
removed = purge_stale_messages_from_queues(
BROKER_URL, (TEST_QUEUE,), business_task_ids=biz_ids, celery_task_ids=celery_ids
)
assert removed == 1
assert client.llen(TEST_QUEUE) == 1 # 正常任务消息保留
client.delete(TEST_QUEUE)