diff --git a/apps/api/app/services/lipsync_service.py b/apps/api/app/services/lipsync_service.py index f0080e972..352657a3b 100644 --- a/apps/api/app/services/lipsync_service.py +++ b/apps/api/app/services/lipsync_service.py @@ -250,7 +250,7 @@ class LipsyncService: gpu_task.id, job.output_duration, ) - # 转存到持久 OSS 路径(GPU 结果已在 gpu-lipsync/results/ 下,直接签短链) + # output_video_url 已是 _submit_to_gpu 内签好的 7 天预签名 URL return # wait_for_result 返回 None 表示超时/最终失败 → 继续走 MediaKit 兜底 logger.warning("[lipsync] GPU 任务等待超时或失败,回退 MediaKit: job_id=%s", job.id) @@ -293,22 +293,77 @@ class LipsyncService: # ── GPU MuseTalk 路径 ──────────────────────────────────────────────── + def _is_own_oss_url(self, url: str, storage) -> bool: + """判断 URL / 存储 key 是否属于自家 OSS。 + + - 裸存储 key(无 scheme):自家对象 + - host 与 storage.public_url host 一致:自家对象 + - 其余 http(s) 公网链接(如 dashscope-result 临时地址):外部对象 + """ + if not url: + return False + parsed = urlparse(url) + if not parsed.scheme: + return True # 裸存储 key + public_base = getattr(storage, "public_url", "") + own_host = urlparse(public_base).netloc.lower() if public_base else "" + return bool(own_host) and parsed.netloc.lower() == own_host + + def _persist_external_audio_for_gpu(self, *, job, storage) -> Optional[str]: + """GPU 任务创建前,把外部域名的预合成 TTS 音频转存到自家 OSS。 + + Worker 部署在用户家庭网络,dashscope-result 等第三方临时 OSS 地址 + 可能无法访问;转存后 gpu_svc 在 poll 时会签自家预签名 URL 给 Worker。 + 已是自家 OSS 对象(含裸 key)直接返回 None(无需转存); + 转存失败返回 None,调用方回退使用原始 URL(最坏情况是 Worker 拉取失败, + 服务端重试耗尽后回退 MediaKit,不阻断业务)。 + """ + if self._is_own_oss_url(job.audio_url, storage): + return None + try: + audio_data = safe_download_bytes( + job.audio_url, + purpose="lipsync_gpu_tts_audio", + allowed_mime_types=ALLOWED_AUDIO_MIME_TYPES, + timeout=60.0, + ) + storage_key = f"lipsync-tts/{job.user_id}/{job.id}.mp3" + permanent_url = storage.upload_file(io.BytesIO(audio_data), storage_key, content_type="audio/mpeg") + logger.info( + "[lipsync] GPU 任务外部音频已转存自家 OSS: job_id=%s key=%s", + job.id, + storage_key, + ) + return permanent_url + except Exception as exc: + logger.warning( + "[lipsync] GPU 任务外部音频转存 OSS 失败,回退原始 URL: job_id=%s err=%s", + job.id, + exc, + ) + return None + def _submit_to_gpu(self, *, job, gpu_svc) -> Optional[object]: """创建 GPU 任务并同步等待结果。 成功返回终态 task 对象(status=done);超时或 GPU 最终失败返回 None, 调用方回退 MediaKit。 - 注意:job.video_url / job.audio_url 可能是: - - 自家 OSS 存储 key(storage.is_own_url 判断,gpu_svc.create_task 内部 - get_download_url 会自动签预签名 URL 给 Worker) - - 外部公网 URL(CosyVoice 临时链接等):poll 返回时原样透传给 Worker, - Worker 可直接 GET 下载。 + 输入处理: + - job.video_url 为用户上传视频,已在自家 OSS(裸 key 或自家 URL), + gpu_svc 在 poll 时签预签名 URL 给 Worker。 + - job.audio_url 可能是预合成 TTS 的第三方临时地址(如 + dashscope-result-bj.oss-cn-beijing.aliyuncs.com),Worker 家庭网络 + 拉不到;创建任务前先转存自家 OSS 再传入。 """ + storage = get_shared_storage_service() + # 外部音频(dashscope 临时链接等)先转存自家 OSS,避免 Worker 家庭网络拉取失败 + persisted_audio_url = self._persist_external_audio_for_gpu(job=job, storage=storage) + audio_url_for_task = persisted_audio_url or job.audio_url # 创建 GPU 任务 gpu_task = gpu_svc.create_task( video_url=job.video_url, - audio_url=job.audio_url, + audio_url=audio_url_for_task, lipsync_job_id=job.id, user_id=job.user_id, project_id=job.project_id, @@ -331,9 +386,21 @@ class LipsyncService: final_task.error_msg, ) return None - # result_url 是 OSS 存储 key;签一个长有效期 URL 写回 job.output_video_url - result_signed = self._sign_media_url(final_task.result_url) - final_task.result_url = result_signed or final_task.result_url + # result_url 是 OSS 存储 key(gpu-lipsync/results/{task_id}.mp4,无 host, + # _sign_media_url 对裸 key 不会签名);直接用 storage 签 7 天预签名 URL + # 写回 job.output_video_url,保证前端拿到可直接下载播放的地址 + try: + signed_result_url = storage.get_download_url( + final_task.result_url, expires_seconds=MEDIAKIT_URL_TTL_SECONDS + ) + if signed_result_url: + final_task.result_url = signed_result_url + except Exception as exc: + logger.warning( + "[lipsync] GPU 结果视频签名失败,回退原始 result_url: gpu_task=%s err=%s", + gpu_task.id, + exc, + ) return final_task # ── 创建任务 ────────────────────────────────────────────────────────── diff --git a/tests/unit/test_lipsync_gpu_integration.py b/tests/unit/test_lipsync_gpu_integration.py index 98c978faa..72a7f2862 100644 --- a/tests/unit/test_lipsync_gpu_integration.py +++ b/tests/unit/test_lipsync_gpu_integration.py @@ -20,7 +20,7 @@ def fake_mediakit(): return client -def _make_job(video_url="oss://video.mp4", audio_url="oss://audio.wav"): +def _make_job(video_url="videos/video.mp4", audio_url="audios/audio.wav"): job = MagicMock() job.id = "job-1" job.user_id = "u1" @@ -42,6 +42,14 @@ def _make_svc(db, mediakit, use_gpu=False): return svc +def _patch_storage(public_url="https://own-bucket.oss-cn-beijing.aliyuncs.com", signed_suffix="?signed-7d"): + """patch get_shared_storage_service,返回自家 OSS storage mock.""" + storage = MagicMock() + storage.public_url = public_url + storage.get_download_url.side_effect = lambda key_or_url, expires_seconds=3600: key_or_url + signed_suffix + return patch("app.services.lipsync_service.get_shared_storage_service", return_value=storage) + + class TestGpuFallback: def test_switch_off_uses_mediakit(self, fake_db, fake_mediakit): """开关关闭时直接走 MediaKit,不调用 _submit_to_gpu.""" @@ -66,26 +74,34 @@ class TestGpuFallback: assert job.status == "submitted" def test_gpu_success_marks_completed(self, fake_db, fake_mediakit): - """GPU 路径成功:job 直接 completed,不调 MediaKit.""" + """GPU 路径成功:job 直接 completed,不调 MediaKit;结果 key 由 storage 签 7 天 URL.""" svc = _make_svc(fake_db, fake_mediakit, use_gpu=True) gpu_done = MagicMock( id="gpu-task-1", status="done", - result_url="oss://gpu-results/r.mp4", + result_url="gpu-lipsync/results/gpu-task-1.mp4", result_duration=12.5, ) fake_gpu_svc = MagicMock() fake_gpu_svc.has_available_worker.return_value = True fake_gpu_svc.create_task.return_value = MagicMock(id="gpu-task-1") fake_gpu_svc.wait_for_result.return_value = gpu_done - with patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc): + with ( + _patch_storage() as storage_p, + patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc), + ): + storage = storage_p() job = _make_job() svc._submit_audio_direct(job=job) fake_gpu_svc.create_task.assert_called_once() fake_mediakit.submit_lipsync.assert_not_called() assert job.status == "completed" assert job.output_duration == 12.5 - assert "?signed" in job.output_video_url + # Bug1 回归:裸 result key 必须经 storage.get_download_url 签 7 天,前端才可播放 + storage.get_download_url.assert_called_once_with( + "gpu-lipsync/results/gpu-task-1.mp4", expires_seconds=7 * 24 * 3600 + ) + assert job.output_video_url == "gpu-lipsync/results/gpu-task-1.mp4?signed-7d" fake_db.commit.assert_called() def test_gpu_timeout_falls_back(self, fake_db, fake_mediakit): @@ -120,12 +136,85 @@ class TestGpuFallback: fake_gpu_svc = MagicMock() fake_gpu_svc.has_available_worker.return_value = True fake_gpu_svc.create_task.side_effect = RuntimeError("DB down") - with patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc): + with _patch_storage(), patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc): job = _make_job() svc._submit_audio_direct(job=job) fake_mediakit.submit_lipsync.assert_called_once() assert job.status == "submitted" + def test_gpu_external_audio_persisted_to_own_oss(self, fake_db, fake_mediakit): + """Bug2 回归:dashscope 临时音频 URL 在创建 GPU 任务前转存自家 OSS。""" + svc = _make_svc(fake_db, fake_mediakit, use_gpu=True) + dashscope_url = "https://dashscope-result-bj.oss-cn-beijing.aliyuncs.com/tmp/abc.mp3" + job = _make_job(audio_url=dashscope_url) + gpu_done = MagicMock(id="gpu-task-2", status="done", result_url="gpu-lipsync/results/gpu-task-2.mp4") + fake_gpu_svc = MagicMock() + fake_gpu_svc.has_available_worker.return_value = True + fake_gpu_svc.create_task.return_value = MagicMock(id="gpu-task-2") + fake_gpu_svc.wait_for_result.return_value = gpu_done + with ( + _patch_storage() as storage_p, + patch("app.services.lipsync_service.safe_download_bytes", return_value=b"FAKE-MP3") as m_dl, + patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc), + ): + storage = storage_p() + storage.upload_file.return_value = "https://own-bucket.oss-cn-beijing.aliyuncs.com/lipsync-tts/u1/job-1.mp3" + svc._submit_audio_direct(job=job) + # 外部音频在 GPU 分支被额外下载(purpose 区分于前置 ffprobe 下载)并转存到约定 key + gpu_dl_calls = [c for c in m_dl.call_args_list if c.kwargs.get("purpose") == "lipsync_gpu_tts_audio"] + assert len(gpu_dl_calls) == 1 + assert gpu_dl_calls[0].args[0] == dashscope_url + storage.upload_file.assert_called_once() + args, kwargs = storage.upload_file.call_args + assert args[1] == "lipsync-tts/u1/job-1.mp3" + assert kwargs.get("content_type") == "audio/mpeg" + # 创建 GPU 任务时用的是自家 OSS URL,Worker 可经预签名下载 + kwargs_create = fake_gpu_svc.create_task.call_args.kwargs + assert kwargs_create["audio_url"] == ("https://own-bucket.oss-cn-beijing.aliyuncs.com/lipsync-tts/u1/job-1.mp3") + assert kwargs_create["audio_url"] != dashscope_url + + def test_gpu_own_audio_not_repersisted(self, fake_db, fake_mediakit): + """Bug2:已是自家 OSS 的音频(含裸 key)不重复下载转存。""" + svc = _make_svc(fake_db, fake_mediakit, use_gpu=True) + job = _make_job(audio_url="lipsync-tts/u1/job-1.mp3") + gpu_done = MagicMock(id="gpu-task-3", status="done", result_url="gpu-lipsync/results/gpu-task-3.mp4") + fake_gpu_svc = MagicMock() + fake_gpu_svc.has_available_worker.return_value = True + fake_gpu_svc.create_task.return_value = MagicMock(id="gpu-task-3") + fake_gpu_svc.wait_for_result.return_value = gpu_done + with ( + _patch_storage() as storage_p, + patch("app.services.lipsync_service.safe_download_bytes") as m_dl, + patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc), + ): + storage = storage_p() + svc._submit_audio_direct(job=job) + # 前置 ffprobe 下载允许发生,但 GPU 转存分支不应再下载/上传 + gpu_dl_calls = [c for c in m_dl.call_args_list if c.kwargs.get("purpose") == "lipsync_gpu_tts_audio"] + assert gpu_dl_calls == [] + storage.upload_file.assert_not_called() + assert fake_gpu_svc.create_task.call_args.kwargs["audio_url"] == "lipsync-tts/u1/job-1.mp3" + + def test_gpu_external_audio_persist_fail_falls_back_original_url(self, fake_db, fake_mediakit): + """Bug2:外部音频转存失败不阻断,用原始 URL 建任务(失败后服务端重试/回退 MediaKit)。""" + svc = _make_svc(fake_db, fake_mediakit, use_gpu=True) + dashscope_url = "https://dashscope-result-bj.oss-cn-beijing.aliyuncs.com/tmp/abc.mp3" + job = _make_job(audio_url=dashscope_url) + gpu_done = MagicMock(id="gpu-task-4", status="done", result_url="gpu-lipsync/results/gpu-task-4.mp4") + fake_gpu_svc = MagicMock() + fake_gpu_svc.has_available_worker.return_value = True + fake_gpu_svc.create_task.return_value = MagicMock(id="gpu-task-4") + fake_gpu_svc.wait_for_result.return_value = gpu_done + with ( + _patch_storage() as storage_p, + patch("app.services.lipsync_service.safe_download_bytes", side_effect=RuntimeError("network blocked")), + patch("app.services.gpu_lipsync_service.GpuLipsyncService", return_value=fake_gpu_svc), + ): + storage = storage_p() + svc._submit_audio_direct(job=job) + storage.upload_file.assert_not_called() + assert fake_gpu_svc.create_task.call_args.kwargs["audio_url"] == dashscope_url + class TestGpuServiceHelpers: """GpuLipsyncService.has_available_worker 测试."""