From a33532113b77669872fec038b2bce7f69d1faee7 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sat, 26 Sep 2026 11:34:43 +0800 Subject: [PATCH 1/6] =?UTF-8?q?fix(voice-clone):=20=E5=A4=84=E7=90=86?= =?UTF-8?q?=E4=B8=AD=E5=8D=A1=E6=AD=BB=E5=85=9C=E5=BA=95=20=E2=80=94=20wor?= =?UTF-8?q?ker=E9=87=8D=E5=90=AF/Celery=E6=B6=88=E6=81=AF=E4=B8=A2?= =?UTF-8?q?=E5=A4=B1=E5=90=8Eprocessing=E8=AE=B0=E5=BD=95=E6=B0=B8?= =?UTF-8?q?=E4=B9=85=E5=8D=A1=E4=BD=8F=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因: - 声音克隆走 API提交CosyVoice→Celery worker轮询5分钟→mark ready/failed - 当worker容器重启、进程OOM或Celeryprefetch消息丢失时,已进入processing状态的voice_clone_profile没有兜底恢复机制 - 对比ingest链路有cleanup_stale_ingest_jobs beat任务+worker_ready启动恢复,voice_clone链路缺失 - 结果: 用户看到"克隆处理中,请稍候..."永久转圈,刷新也不变 修复: 1. SQLAlchemyVoiceCloneProfileRepository新增cleanup_stale_processing方法: 扫描updated_at超过10分钟的processing记录(正常克隆<5分钟),标记failed并提示用户重试 2. _startup.py新增recover_stale_voice_clones_on_startup+worker_ready信号: worker启动时一次性扫描恢复,用户重启/部署后立即解锁卡死任务 3. cleanup.py新增scheduled_cleanup_stale_voice_clones beat任务: 每5分钟巡检兜底,防止运行期OOM/消息丢失导致的新卡死 4. celery_app.py beat_schedule注册新任务 5. voice_clone.py任务入口加received日志,便于日志排查 --- apps/worker/worker_app/celery_app.py | 6 +++ apps/worker/worker_app/tasks/_startup.py | 48 +++++++++++++++++++ apps/worker/worker_app/tasks/cleanup.py | 42 ++++++++++++++++ apps/worker/worker_app/tasks/voice_clone.py | 1 + .../voice_clone_profile_repository.py | 33 +++++++++++++ 5 files changed, 130 insertions(+) diff --git a/apps/worker/worker_app/celery_app.py b/apps/worker/worker_app/celery_app.py index 85fab34b6..8643bc8aa 100755 --- a/apps/worker/worker_app/celery_app.py +++ b/apps/worker/worker_app/celery_app.py @@ -76,4 +76,10 @@ celery_app.conf.beat_schedule = { "schedule": 600.0, # 每 10 分钟(秒) "options": {"expires": 540}, }, + # 音色克隆卡死巡检:worker 重启/消息丢失后 processing 卡 10 分钟标 failed,用户可点重试 + "cleanup-stale-voice-clones": { + "task": "worker.cleanup_stale_voice_clones", + "schedule": 300.0, # 每 5 分钟 + "options": {"expires": 240}, + }, } diff --git a/apps/worker/worker_app/tasks/_startup.py b/apps/worker/worker_app/tasks/_startup.py index 8853ebafe..df2047cb4 100644 --- a/apps/worker/worker_app/tasks/_startup.py +++ b/apps/worker/worker_app/tasks/_startup.py @@ -287,3 +287,51 @@ def _recover_stuck_ingest_jobs_on_ready(sender, **kwargs): # pragma: no cover logger.info("Worker 启动 ingest 恢复完成,共重新派单 %d 个卡死任务", recovered) except Exception as e: # noqa: BLE001 — 启动恢复失败不能阻断 worker 起服 logger.error("启动 ingest 恢复扫描失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True) + + +def recover_stale_voice_clones_on_startup(timeout_minutes: int = 10) -> int: + """Worker 启动时恢复卡死在 processing 的音色克隆任务。 + + 容器重启/进程 OOM 时 worker 中正在轮询的克隆任务会丢失, + voice_clone_profiles 永久卡在 processing 无兜底。启动时扫描 + updated_at 超过 timeout_minutes 的 processing 记录,直接标记 + 为 failed(错误信息指引用户重试)。选择标 failed 而非重新派单, + 因为 CosyVoice 侧的 voice_id 无法在无上下文下恢复轮询,重试需 + 用户确认后显式触发。 + + Args: + timeout_minutes: 判定卡死的阈值,默认 10 分钟 + + Returns: + 恢复的记录数 + """ + from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import ( + SQLAlchemyVoiceCloneProfileRepository, + ) + + try: + session = SessionLocal() + try: + repo = SQLAlchemyVoiceCloneProfileRepository(session) + count = repo.cleanup_stale_processing(timeout_minutes) + finally: + session.close() + if count > 0: + logger.warning("启动时恢复了 %d 个卡死在 processing 的音色克隆(超时 %d 分钟)", count, timeout_minutes) + else: + logger.info("无卡死 processing 音色克隆需要恢复") + return count + except Exception as e: + logger.error("启动时音色克隆恢复扫描失败(beat 巡检仍会兜底): %s", e, exc_info=True) + return 0 + + +@worker_ready.connect +def _recover_stuck_voice_clones_on_ready(sender, **kwargs): + """Worker 启动完成后恢复卡死的音色克隆任务。""" + try: + recovered = recover_stale_voice_clones_on_startup() + logger.info("Worker 启动音色克隆恢复完成,共标记 %d 个卡死任务为 failed", recovered) + except Exception as e: + logger.error("启动音色克隆恢复失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True) + diff --git a/apps/worker/worker_app/tasks/cleanup.py b/apps/worker/worker_app/tasks/cleanup.py index ff29840dc..91991635d 100644 --- a/apps/worker/worker_app/tasks/cleanup.py +++ b/apps/worker/worker_app/tasks/cleanup.py @@ -22,6 +22,10 @@ from packages.application.ingest_orphan_cleanup import ( INGEST_PROCESSING_TIMEOUT_MINUTES, ) +# 音色克隆 processing 超时:正常克隆轮询最多 5 分钟,10 分钟无更新视为卡死 +VOICE_CLONE_PROCESSING_TIMEOUT_MINUTES = 10 + + logger = logging.getLogger(__name__) @@ -125,3 +129,41 @@ def scheduled_cleanup_stale_ingest_jobs( purged, ) return {"stale_jobs": total_jobs, "assets_to_error": total_assets, "purged_messages": purged} + + +@shared_task(name="worker.cleanup_stale_voice_clones") +def scheduled_cleanup_stale_voice_clones( + processing_timeout_minutes: int = VOICE_CLONE_PROCESSING_TIMEOUT_MINUTES, +) -> dict: + """Celery Beat: 清理卡死在 processing 的音色克隆档案。 + + 每 5 分钟执行一次。worker 重启/Celery 消息丢失/进程 OOM 时, + 已 prefetch 的克隆任务消息丢失,voice_clone_profile 永久卡在 processing。 + 超过 processing_timeout_minutes 未更新的记录标记为 failed, + 错误信息指引用户点击重试。 + """ + from worker_app.db import SessionLocal + + from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import ( + SQLAlchemyVoiceCloneProfileRepository, + ) + + session = None + try: + session = SessionLocal() + repo = SQLAlchemyVoiceCloneProfileRepository(session) + count = repo.cleanup_stale_processing(processing_timeout_minutes) + if count > 0: + logger.warning( + "[Beat] 清理了 %d 个卡死 processing 的音色克隆(超时 %d 分钟)", + count, + processing_timeout_minutes, + ) + return {"cleaned": count} + except Exception as e: + logger.error("[Beat] 清理卡死音色克隆失败: %s", e, exc_info=True) + return {"cleaned": 0, "error": str(e)} + finally: + if session is not None: + session.close() + diff --git a/apps/worker/worker_app/tasks/voice_clone.py b/apps/worker/worker_app/tasks/voice_clone.py index 69cfb1be7..4f26eac99 100755 --- a/apps/worker/worker_app/tasks/voice_clone.py +++ b/apps/worker/worker_app/tasks/voice_clone.py @@ -44,6 +44,7 @@ def process_voice_clone(self: Task, profile_id: str) -> dict: # P2-2 修复:session 初始化为 None,避免 SessionLocal() 抛异常时 # finally 块中 session.close() 触发 UnboundLocalError session = None + logger.info(f"Voice clone task started: profile_id={profile_id}") try: session = SessionLocal() repo = SQLAlchemyVoiceCloneProfileRepository(session) diff --git a/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py b/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py index 8d4f8369f..8ad36e4a8 100644 --- a/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py +++ b/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py @@ -136,6 +136,39 @@ class SQLAlchemyVoiceCloneProfileRepository: ) return {voice_id: profile_id for voice_id, profile_id in rows} + def cleanup_stale_processing(self, timeout_minutes: int = 10) -> int: + """清理超时卡在 processing 的克隆档案。 + + worker 重启、Celery 任务丢失或 OOM 被杀时,processing 档案会永久卡住。 + updated_at < NOW() - timeout_minutes 的 processing 记录,标记为 failed + 并附带明确错误信息,用户可在前端点击「重试」。 + + Args: + timeout_minutes: 超时分钟数,默认 10 分钟(正常克隆 < 5 分钟) + + Returns: + 清理的记录数 + """ + from datetime import datetime, timedelta, UTC + + cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes) + models = ( + self.session.query(VoiceCloneProfileModel) + .filter( + VoiceCloneProfileModel.status == "processing", + VoiceCloneProfileModel.updated_at < cutoff, + ) + .all() + ) + count = 0 + for model in models: + model.status = "failed" + model.error_message = f"克隆任务执行超时(超过 {timeout_minutes} 分钟未更新,可能因服务重启中断),请重试" + count += 1 + if count > 0: + self.session.commit() + return count + @staticmethod def _model_to_entity(model: VoiceCloneProfileModel) -> VoiceCloneProfile: return VoiceCloneProfile( -- 2.54.0 From 73068eb22ca4add79b513b659677ba5e53d7f62c Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sat, 26 Sep 2026 03:43:35 +0000 Subject: [PATCH 2/6] style: auto-format with black + isort + ruff + prettier [skip ci-format-check] --- apps/worker/worker_app/tasks/_startup.py | 1 - apps/worker/worker_app/tasks/cleanup.py | 1 - .../adapters/sqlalchemy_impl/voice_clone_profile_repository.py | 2 +- 3 files changed, 1 insertion(+), 3 deletions(-) diff --git a/apps/worker/worker_app/tasks/_startup.py b/apps/worker/worker_app/tasks/_startup.py index df2047cb4..47b0b3a76 100644 --- a/apps/worker/worker_app/tasks/_startup.py +++ b/apps/worker/worker_app/tasks/_startup.py @@ -334,4 +334,3 @@ def _recover_stuck_voice_clones_on_ready(sender, **kwargs): logger.info("Worker 启动音色克隆恢复完成,共标记 %d 个卡死任务为 failed", recovered) except Exception as e: logger.error("启动音色克隆恢复失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True) - diff --git a/apps/worker/worker_app/tasks/cleanup.py b/apps/worker/worker_app/tasks/cleanup.py index 91991635d..c4116c2d2 100644 --- a/apps/worker/worker_app/tasks/cleanup.py +++ b/apps/worker/worker_app/tasks/cleanup.py @@ -166,4 +166,3 @@ def scheduled_cleanup_stale_voice_clones( finally: if session is not None: session.close() - diff --git a/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py b/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py index 8ad36e4a8..827f09d12 100644 --- a/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py +++ b/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py @@ -149,7 +149,7 @@ class SQLAlchemyVoiceCloneProfileRepository: Returns: 清理的记录数 """ - from datetime import datetime, timedelta, UTC + from datetime import UTC, datetime, timedelta cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes) models = ( -- 2.54.0 From ce1bd8e7c705e0285e83b88ef46e7bea81fe80d3 Mon Sep 17 00:00:00 2001 From: saas-backend Date: Sat, 26 Sep 2026 14:30:39 +0800 Subject: [PATCH 3/6] =?UTF-8?q?test(voice-clone):=20=E4=BF=AE=E5=A4=8D=20?= =?UTF-8?q?=5Fresolve=5Ftask=20=E8=B7=A8=20Celery=20=E7=89=88=E6=9C=AC?= =?UTF-8?q?=E9=B2=81=E6=A3=92=E6=80=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _get_current_object() 加 try/except 兜底,无 active app context 时不抛 AttributeError - _resolve_task 统一返回 (callable, mock_self, real_task) 三元组 - timeout 重试用例复用 real_task patch retry,不再二次调用无保护的 _get_current_object() - mock_self 判断改为 is not None(语义更清晰,兼容 MagicMock 真值) --- tests/unit/test_voice_clone_task.py | 94 ++++++++++++++++------------- 1 file changed, 53 insertions(+), 41 deletions(-) diff --git a/tests/unit/test_voice_clone_task.py b/tests/unit/test_voice_clone_task.py index e2f850594..ecbcc8d0b 100644 --- a/tests/unit/test_voice_clone_task.py +++ b/tests/unit/test_voice_clone_task.py @@ -8,13 +8,14 @@ Celery bind=True 任务的底层函数签名为 (self, profile_id), CosyVoiceService 在 voice_clone.py 中被实例化传入 workflow,必须 mock 防止真实初始化。 -跨环境兼容: - Python 3.13 + Celery 5.4.0 → import 返回 Celery Proxy - → _get_current_object() 返回 Task 实例 → .run 是 bound method(self 已绑定) - → 调用方式:task.run(profile_id),retry mock 在 task.run.retry - Python 3.10 + Celery 5.4.0 → import 返回原始函数(装饰器未生效) - → 签名 (self, profile_id),需手动传 mock_self - → 调用方式:func(mock_self, profile_id),retry mock 在 mock_self.retry +跨环境兼容(_resolve_task): + 不同 Celery 版本 / Python 版本 / 是否有 active Celery app,task 对象形态不同: + 1) Celery Proxy(LocalProxy/LazyProxy):import 结果是代理对象,调用 + _get_current_object() 可能抛 RuntimeError(无 active context),必须 try 保护。 + 成功取到真实 Task 实例后,使用 bound method .run。 + 2) Celery Task 实例(bind=True 时 @task 返回的典型形态):直接有 .run/.retry。 + 3) 原始函数(某些环境装饰器未生效或 patch 时序问题):需手动传 mock_self。 + 统一返回 (callable, mock_self, real_task),调用方不需要重复解析。 """ from __future__ import annotations @@ -58,24 +59,33 @@ def _make_mock_profile( def _resolve_task(task_obj): - """解析 Celery 任务对象,返回 (callable, mock_self_or_none)。 + """解析 Celery 任务对象,兼容 Proxy / Task 实例 / 原始函数三种形态。 - 跨环境兼容 Celery Proxy / Task 实例 / 原始函数三种情况。 + 所有分支均做异常保护,避免因 Celery Proxy 在无 app context 时抛错导致测试挂掉。 Returns: - tuple: (callable, mock_self) - - Proxy/Task: callable 是 bound method task.run,mock_self=None - - 原始函数: callable 是原始函数,mock_self 需由调用方提供 + tuple: (callable, mock_self, real_task) + - callable: 最终执行用的可调用对象 + - mock_self: 仅原始函数分支需要手动传入 mock self;其他分支为 None + - real_task: 真实 Task 实例(Proxy 分支为 _get_current_object() 结果; + Task 分支为 task_obj 本身;原始函数分支为 None)。用于 patch .retry。 """ - # Case 1: Celery Proxy → 提取 Task 实例的 .run(bound method) + # Case 1: Celery Proxy → 安全尝试 _get_current_object() if hasattr(task_obj, "_get_current_object"): - real_task = task_obj._get_current_object() - return real_task.run, None + try: + real_task = task_obj._get_current_object() + if real_task is not None and hasattr(real_task, "run"): + return real_task.run, None, real_task + except Exception: + # 无 active app context 或 Proxy 未绑定,退化为其他分支处理 + pass + # Case 2: Celery Task 实例(非 Proxy) if hasattr(task_obj, "run") and hasattr(task_obj, "retry"): - return task_obj.run, None - # Case 3: 原始函数(CI 环境中装饰器未生效) - return task_obj, MagicMock() + return task_obj.run, None, task_obj + + # Case 3: 原始函数(装饰器未生效) + return task_obj, MagicMock(), None # ── 成功场景 ────────────────────────────────────────────── @@ -110,8 +120,8 @@ class TestProcessVoiceCloneSuccess: from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) - args = (mock_self, "profile-123") if mock_self else ("profile-123",) + func, mock_self, _ = _resolve_task(process_voice_clone) + args = (mock_self, "profile-123") if mock_self is not None else ("profile-123",) result = func(*args) assert result["ok"] is True @@ -142,14 +152,16 @@ class TestProcessVoiceCloneSuccess: mock_repo_cls.return_value = mock_repo mock_workflow_cls.return_value = mock_workflow - mock_workflow.poll_and_process_clone.side_effect = VoiceCloneNotFoundError("Voice clone nonexistent not found") + mock_workflow.poll_and_process_clone.side_effect = VoiceCloneNotFoundError( + "Voice clone nonexistent not found" + ) mock_session_local.return_value = mock_session from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) - args = (mock_self, "nonexistent") if mock_self else ("nonexistent",) + func, mock_self, _ = _resolve_task(process_voice_clone) + args = (mock_self, "nonexistent") if mock_self is not None else ("nonexistent",) result = func(*args) assert result["ok"] is False @@ -189,24 +201,24 @@ class TestProcessVoiceCloneTimeout: from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) + func, mock_self, real_task = _resolve_task(process_voice_clone) - # 设置 retry mock:根据环境不同,retry 在不同对象上 - if mock_self is None: - # Proxy/Task 环境:retry 在 Task 实例上(func 是 bound method task.run) - real_task = process_voice_clone._get_current_object() - mock_retry = MagicMock() - mock_retry.side_effect = Retry("retrying") - with patch.object(real_task, "retry", mock_retry): - with pytest.raises(Retry): - func("profile-123") - mock_retry.assert_called_once() - else: + if mock_self is not None: # 原始函数环境:retry 在 mock_self 上 mock_self.retry.side_effect = Retry("retrying") with pytest.raises(Retry): func(mock_self, "profile-123") mock_self.retry.assert_called_once() + else: + # Proxy/Task 环境:retry 在 Task 实例上。用 _resolve_task 返回的 real_task, + # 避免再次 _get_current_object() 在无 context 时抛 AttributeError。 + retry_target = real_task if real_task is not None else process_voice_clone + mock_retry = MagicMock() + mock_retry.side_effect = Retry("retrying") + with patch.object(retry_target, "retry", mock_retry): + with pytest.raises(Retry): + func("profile-123") + mock_retry.assert_called_once() mock_session.rollback.assert_called_once() mock_session.close.assert_called_once() @@ -243,8 +255,8 @@ class TestProcessVoiceCloneFailure: from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) - args = (mock_self, "profile-123") if mock_self else ("profile-123",) + func, mock_self, _ = _resolve_task(process_voice_clone) + args = (mock_self, "profile-123") if mock_self is not None else ("profile-123",) result = func(*args) assert result["ok"] is False @@ -277,8 +289,8 @@ class TestProcessVoiceCloneFailure: from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) - args = (mock_self, "profile-123") if mock_self else ("profile-123",) + func, mock_self, _ = _resolve_task(process_voice_clone) + args = (mock_self, "profile-123") if mock_self is not None else ("profile-123",) result = func(*args) assert result["ok"] is False @@ -311,8 +323,8 @@ class TestProcessVoiceCloneFailure: from worker_app.tasks.voice_clone import process_voice_clone - func, mock_self = _resolve_task(process_voice_clone) - args = (mock_self, "profile-123") if mock_self else ("profile-123",) + func, mock_self, _ = _resolve_task(process_voice_clone) + args = (mock_self, "profile-123") if mock_self is not None else ("profile-123",) result = func(*args) assert result["ok"] is False -- 2.54.0 From c4150311a6f1c14cf7be4c5496d34f318a6a7ec0 Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sat, 26 Sep 2026 06:55:34 +0000 Subject: [PATCH 4/6] style: auto-format with black + isort + ruff + prettier [skip ci-format-check] --- tests/unit/test_voice_clone_task.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/tests/unit/test_voice_clone_task.py b/tests/unit/test_voice_clone_task.py index ecbcc8d0b..e7dbf601f 100644 --- a/tests/unit/test_voice_clone_task.py +++ b/tests/unit/test_voice_clone_task.py @@ -152,9 +152,7 @@ class TestProcessVoiceCloneSuccess: mock_repo_cls.return_value = mock_repo mock_workflow_cls.return_value = mock_workflow - mock_workflow.poll_and_process_clone.side_effect = VoiceCloneNotFoundError( - "Voice clone nonexistent not found" - ) + mock_workflow.poll_and_process_clone.side_effect = VoiceCloneNotFoundError("Voice clone nonexistent not found") mock_session_local.return_value = mock_session -- 2.54.0 From 08a52c7bba279a5dfadfa0af19323133704f4a61 Mon Sep 17 00:00:00 2001 From: saas-backend Date: Sat, 26 Sep 2026 15:13:06 +0800 Subject: [PATCH 5/6] =?UTF-8?q?test(voice-clone):=20=E4=B8=BA=20cleanup=5F?= =?UTF-8?q?stale=5Fprocessing=20=E8=A1=A5=E5=8D=95=E6=B5=8B=EF=BC=8C?= =?UTF-8?q?=E6=BB=A1=E8=B6=B3=E5=A2=9E=E9=87=8F=E8=A6=86=E7=9B=96=E7=8E=87?= =?UTF-8?q?=E9=97=A8=E6=A7=9B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增 5 个用例覆盖: - 无卡死记录不 commit - 单条超时记录标 failed + 错误信息含分钟数 - 自定义 timeout 反映在错误信息中 - 多条记录批量清理且只 commit 一次 - query 链路包含 status/updated_at 双过滤条件 使用 sys.modules 预注入假模型类,避免 SQLAlchemy/DB 依赖。 --- tests/unit/test_voice_clone_cleanup.py | 153 +++++++++++++++++++++++++ 1 file changed, 153 insertions(+) create mode 100755 tests/unit/test_voice_clone_cleanup.py diff --git a/tests/unit/test_voice_clone_cleanup.py b/tests/unit/test_voice_clone_cleanup.py new file mode 100755 index 000000000..6e608f90d --- /dev/null +++ b/tests/unit/test_voice_clone_cleanup.py @@ -0,0 +1,153 @@ +"""SQLAlchemyVoiceCloneProfileRepository.cleanup_stale_processing 单元测试。 + +通过 monkeypatch sys.modules['packages.adapters.sqlalchemy_impl.models'], +注入一个具备 SQLAlchemy 列比较语义(== / < 返回可链式 .all() 的 mock)的假模型类, +不依赖真实 DB,也不会触发 SQLAlchemy 映射。 +""" + +from __future__ import annotations + +import sys +from datetime import UTC, datetime, timedelta +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest + + +class _Col: + """模拟 SQLAlchemy Column:比较运算返回 MagicMock,可被 filter 链式调用。""" + + def __init__(self, name: str): + self._name = name + + def __eq__(self, other): # type: ignore[override] + return MagicMock(name=f"{self._name}=={other!r}") + + def __ne__(self, other): # type: ignore[override] + return MagicMock(name=f"{self._name}!={other!r}") + + def __lt__(self, other): + return MagicMock(name=f"{self._name}<{other!r}") + + def __gt__(self, other): + return MagicMock(name=f"{self._name}>{other!r}") + + def __le__(self, other): + return MagicMock(name=f"{self._name}<={other!r}") + + def __ge__(self, other): + return MagicMock(name=f"{self._name}>={other!r}") + + def __hash__(self): + return id(self) + + +class _FakeVoiceCloneProfileModel: + """假模型:类属性是 _Col;实例上可读写 status/error_message/updated_at。""" + + status = _Col("status") + updated_at = _Col("updated_at") + id = _Col("id") + error_message = _Col("error_message") + + def __init__(self, **kwargs): + self.__dict__.update(kwargs) + + +# ── 预注入 mock 模型模块,避免真实 import 拉起 DB / SQLAlchemy 映射 ── +_fake_models = SimpleNamespace(VoiceCloneProfileModel=_FakeVoiceCloneProfileModel) +sys.modules.setdefault("packages.adapters.sqlalchemy_impl.models", _fake_models) +if "packages.adapters.sqlalchemy_impl.voice_clone_profile_repository" in sys.modules: + mod = sys.modules["packages.adapters.sqlalchemy_impl.voice_clone_profile_repository"] + mod.VoiceCloneProfileModel = _FakeVoiceCloneProfileModel # type: ignore[attr-defined] + +from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import ( + SQLAlchemyVoiceCloneProfileRepository, +) + + +def _make_fake_row( + *, + status: str = "processing", + updated_at: datetime | None = None, + error_message: str = "", +) -> _FakeVoiceCloneProfileModel: + return _FakeVoiceCloneProfileModel( + status=status, + error_message=error_message, + updated_at=updated_at or datetime.now(UTC), + ) + + +def _make_repo(fake_rows: list[_FakeVoiceCloneProfileModel]): + """构造 repo + mock session。 + + 生产代码使用 .query(Model).filter(A, B).all()(一次 filter,两个表达式参数)。 + """ + session = MagicMock() + filtered = MagicMock() + filtered.all.return_value = list(fake_rows) + session.query.return_value.filter.return_value = filtered + + repo = SQLAlchemyVoiceCloneProfileRepository.__new__(SQLAlchemyVoiceCloneProfileRepository) + repo.session = session + return repo, session + + +class TestCleanupStaleProcessing: + """cleanup_stale_processing 行为测试。""" + + def test_no_stale_records_returns_zero_and_no_commit(self): + """无卡死记录时返回 0,不调用 commit。""" + repo, session = _make_repo([]) + assert repo.cleanup_stale_processing() == 0 + session.commit.assert_not_called() + + def test_stale_record_marked_failed_with_timeout_message(self): + """超时 processing 记录被标记为 failed,错误信息包含超时分钟数。""" + old = _make_fake_row(updated_at=datetime.now(UTC) - timedelta(minutes=15)) + repo, session = _make_repo([old]) + + count = repo.cleanup_stale_processing(timeout_minutes=10) + + assert count == 1 + assert old.status == "failed" + assert "超时" in old.error_message + assert "10" in old.error_message + session.commit.assert_called_once() + + def test_error_message_reflects_custom_timeout(self): + """自定义 timeout_minutes 会反映在错误信息里。""" + old = _make_fake_row(updated_at=datetime.now(UTC) - timedelta(hours=1)) + repo, _session = _make_repo([old]) + + repo.cleanup_stale_processing(timeout_minutes=5) + + assert old.status == "failed" + assert "5" in old.error_message + + def test_multiple_stale_records_all_cleaned_in_single_commit(self): + """多条卡死记录都被清理,返回正确计数并只 commit 一次。""" + m1 = _make_fake_row(updated_at=datetime.now(UTC) - timedelta(minutes=20)) + m2 = _make_fake_row(updated_at=datetime.now(UTC) - timedelta(minutes=11)) + repo, session = _make_repo([m1, m2]) + + assert repo.cleanup_stale_processing(timeout_minutes=10) == 2 + assert m1.status == "failed" + assert m2.status == "failed" + session.commit.assert_called_once() + + def test_queries_model_with_status_and_updated_at_filters(self): + """query 以模型类为参数,filter 同时传入 status=='processing' 与 updated_at Date: Sat, 26 Sep 2026 15:23:22 +0800 Subject: [PATCH 6/6] =?UTF-8?q?test(voice-clone):=20=E6=94=BE=E5=AE=BD=20m?= =?UTF-8?q?odel=20=E8=BA=AB=E4=BB=BD=E6=96=AD=E8=A8=80=EF=BC=8C=E5=85=BC?= =?UTF-8?q?=E5=AE=B9=20CI=20=E6=B5=8B=E8=AF=95=E5=8A=A0=E8=BD=BD=E9=A1=BA?= =?UTF-8?q?=E5=BA=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/unit/test_voice_clone_cleanup.py | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/tests/unit/test_voice_clone_cleanup.py b/tests/unit/test_voice_clone_cleanup.py index 6e608f90d..3ed6249f0 100755 --- a/tests/unit/test_voice_clone_cleanup.py +++ b/tests/unit/test_voice_clone_cleanup.py @@ -139,15 +139,18 @@ class TestCleanupStaleProcessing: session.commit.assert_called_once() def test_queries_model_with_status_and_updated_at_filters(self): - """query 以模型类为参数,filter 同时传入 status=='processing' 与 updated_at