Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d5a4a40f2a |
@@ -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
|
||||
|
||||
# ── 创建任务 ──────────────────────────────────────────────────────────
|
||||
|
||||
@@ -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 测试."""
|
||||
|
||||
Reference in New Issue
Block a user