From 02bc1129b3ec451a3d6d56f75546ef446094373a Mon Sep 17 00:00:00 2001 From: CI Bot Date: Sat, 18 Jul 2026 00:17:49 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20ingest=E5=85=88=E4=B8=8B?= =?UTF-8?q?=E8=BD=BDOSS=E6=96=87=E4=BB=B6=E5=86=8D=E6=8F=90=E5=8F=96?= =?UTF-8?q?=E5=85=83=E6=95=B0=E6=8D=AE=EF=BC=8C=E4=BF=AE=E5=A4=8D=E6=9C=89?= =?UTF-8?q?=E6=95=88=E6=80=A7=E6=A0=A1=E9=AA=8C=E5=85=A8=E9=83=A8=E8=AF=AF?= =?UTF-8?q?=E5=88=A4ERROR?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因:#488新增的有效性校验直接把storage_key传给ffprobe/Pillow, 但storage_key是OSS内部路径不是可访问的URL,元数据提取永远失败。 修复: - ingest_asset先调用download_asset下载到临时文件 - 对本地文件跑ffprobe/Pillow提取元数据 - 用完删除临时文件 - 下载失败也标记为ERROR(素材本身不可用) 新增2个测试:下载失败场景 + 临时文件清理验证 --- apps/worker/worker_app/tasks/ingest.py | 26 +++- tests/unit/test_ingest_validation.py | 169 +++++++++++++++++++++---- 2 files changed, 165 insertions(+), 30 deletions(-) mode change 100755 => 100644 apps/worker/worker_app/tasks/ingest.py diff --git a/apps/worker/worker_app/tasks/ingest.py b/apps/worker/worker_app/tasks/ingest.py old mode 100755 new mode 100644 index 47879b913..12941d576 --- a/apps/worker/worker_app/tasks/ingest.py +++ b/apps/worker/worker_app/tasks/ingest.py @@ -1,7 +1,10 @@ import subprocess +import tempfile from datetime import datetime, timezone +from pathlib import Path from celery.utils.log import get_task_logger +from video_processing.oss_helpers import download_asset from worker_app.celery_app import celery_app from worker_app.core.asset_types import infer_mime_type_from_storage_key from worker_app.db import SessionLocal @@ -227,9 +230,26 @@ def ingest_asset(job_id: str) -> dict: elif mime_type.startswith("audio/"): media_type = "audio" - # Extract metadata - storage_url = job.storage_key # Assuming storage_key is usable as URL/path - metadata, extract_success = extract_media_metadata(storage_url, media_type) + # 先从 OSS 下载文件到本地临时目录,再提取元数据 + # (storage_key 是 OSS 内部路径,不能直接传给 ffprobe/Pillow) + local_file = None + try: + suffix = Path(job.storage_key).suffix or ".bin" + with tempfile.NamedTemporaryFile(suffix=suffix, delete=False) as tmp: + local_file = Path(tmp.name) + + download_ok = download_asset(job.storage_key, local_file) + if not download_ok: + logger.warning("素材下载失败,无法提取元数据: job_id=%s storage_key=%s", job_id, job.storage_key) + metadata, extract_success = {}, False + else: + metadata, extract_success = extract_media_metadata(str(local_file), media_type) + finally: + if local_file and local_file.exists(): + try: + local_file.unlink() + except OSError: + pass # 有效性校验:ffprobe/Pillow 必须成功,且文件大小/时长/尺寸满足最小要求 is_valid = extract_success and _is_valid_media(metadata, media_type) diff --git a/tests/unit/test_ingest_validation.py b/tests/unit/test_ingest_validation.py index 31877e6b0..d11f4ada8 100755 --- a/tests/unit/test_ingest_validation.py +++ b/tests/unit/test_ingest_validation.py @@ -171,18 +171,19 @@ class TestIngestAssetValidation: job_repo.create(job) db.commit() - # mock extract_media_metadata 返回失败 - with patch("worker_app.tasks.ingest.extract_media_metadata") as mock_extract: - mock_extract.return_value = ( - {"size_bytes": 27, "duration": 0, "width": 0, "height": 0}, - False, - ) - # mock SessionLocal 返回我们的db session - with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): - # mock celery_app.task decorator不影响函数本身 - from worker_app.tasks.ingest import ingest_asset + # mock download_asset 成功 + extract_media_metadata 返回失败 + with patch("worker_app.tasks.ingest.download_asset", return_value=True): + with patch("worker_app.tasks.ingest.extract_media_metadata") as mock_extract: + mock_extract.return_value = ( + {"size_bytes": 27, "duration": 0, "width": 0, "height": 0}, + False, + ) + # mock SessionLocal 返回我们的db session + with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): + # mock celery_app.task decorator不影响函数本身 + from worker_app.tasks.ingest import ingest_asset - result = ingest_asset("job-invalid-001") + result = ingest_asset("job-invalid-001") assert result["status"] == "failed" assert "asset_id" in result @@ -233,23 +234,24 @@ class TestIngestAssetValidation: job_repo.create(job) db.commit() - with patch("worker_app.tasks.ingest.extract_media_metadata") as mock_extract: - mock_extract.return_value = ( - { - "size_bytes": 5 * 1024 * 1024, - "duration": 30.5, - "width": 1920, - "height": 1080, - "codec": "h264", - "fps": 30.0, - "bitrate": 2000000, - }, - True, - ) - with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): - from worker_app.tasks.ingest import ingest_asset + with patch("worker_app.tasks.ingest.download_asset", return_value=True): + with patch("worker_app.tasks.ingest.extract_media_metadata") as mock_extract: + mock_extract.return_value = ( + { + "size_bytes": 5 * 1024 * 1024, + "duration": 30.5, + "width": 1920, + "height": 1080, + "codec": "h264", + "fps": 30.0, + "bitrate": 2000000, + }, + True, + ) + with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): + from worker_app.tasks.ingest import ingest_asset - result = ingest_asset("job-valid-001") + result = ingest_asset("job-valid-001") assert result["status"] == "completed" assert "asset_id" in result @@ -265,3 +267,116 @@ class TestIngestAssetValidation: assert updated_job.status == IngestJobStatus.COMPLETED db.close() + + def test_download_failure_marked_as_error(self): + """OSS下载失败 → 无法提取元数据 → asset.status=ERROR.""" + from sqlalchemy import create_engine + from sqlalchemy.orm import sessionmaker + + from packages.adapters.sqlalchemy_impl.models import Base + from packages.domain import AssetStatus, IngestJobStatus + + engine = create_engine("sqlite:///:memory:") + Base.metadata.create_all(engine) + Session = sessionmaker(bind=engine) + db = Session() + + from packages.adapters.sqlalchemy_impl import ( + SQLAlchemyAssetRepository, + SQLAlchemyIngestJobRepository, + ) + from packages.domain import IngestJob + + job_repo = SQLAlchemyIngestJobRepository(db) + asset_repo = SQLAlchemyAssetRepository(db) + + job = IngestJob( + id="job-dl-fail-001", + project_id="proj-1", + library_id="lib-1", + storage_key="uploads/test/nonexistent.mp4", + status=IngestJobStatus.PENDING, + ) + job_repo.create(job) + db.commit() + + # mock download_asset 返回失败 + with patch("worker_app.tasks.ingest.download_asset", return_value=False): + with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): + from worker_app.tasks.ingest import ingest_asset + + result = ingest_asset("job-dl-fail-001") + + assert result["status"] == "failed" + assert "asset_id" in result + + asset = asset_repo.get(result["asset_id"]) + assert asset is not None + assert asset.status == AssetStatus.ERROR + + updated_job = job_repo.get("job-dl-fail-001") + assert updated_job.status == IngestJobStatus.FAILED + + db.close() + + def test_temp_file_cleaned_up_after_extraction(self): + """提取元数据后临时文件被清理.""" + import tempfile + + from sqlalchemy import create_engine + from sqlalchemy.orm import sessionmaker + + from packages.adapters.sqlalchemy_impl.models import Base + from packages.domain import IngestJobStatus + + engine = create_engine("sqlite:///:memory:") + Base.metadata.create_all(engine) + Session = sessionmaker(bind=engine) + db = Session() + + from packages.adapters.sqlalchemy_impl import SQLAlchemyIngestJobRepository + from packages.domain import IngestJob + + job_repo = SQLAlchemyIngestJobRepository(db) + + job = IngestJob( + id="job-cleanup-001", + project_id="proj-1", + library_id="lib-1", + storage_key="uploads/test/cleanup.mp4", + status=IngestJobStatus.PENDING, + ) + job_repo.create(job) + db.commit() + + temp_files_created = [] + original_namedtempfile = tempfile.NamedTemporaryFile + + def tracking_namedtempfile(*args, **kwargs): + tmp = original_namedtempfile(*args, **kwargs) + temp_files_created.append(tmp.name) + return tmp + + def fake_download(storage_key, local_path): + # 模拟下载:写点假数据 + Path(local_path).write_bytes(b"fake video data" * 100) + return True + + with patch("tempfile.NamedTemporaryFile", side_effect=tracking_namedtempfile): + with patch("worker_app.tasks.ingest.download_asset", side_effect=fake_download): + with patch("worker_app.tasks.ingest.extract_media_metadata") as mock_extract: + mock_extract.return_value = ( + {"size_bytes": 5 * 1024 * 1024, "duration": 30.5, "width": 1920, "height": 1080}, + True, + ) + with patch("worker_app.tasks.ingest.SessionLocal", return_value=db): + from worker_app.tasks.ingest import ingest_asset + + ingest_asset("job-cleanup-001") + + # 验证临时文件已被清理 + assert len(temp_files_created) > 0 + for f in temp_files_created: + assert not Path(f).exists(), f"临时文件未被清理: {f}" + + db.close() -- 2.54.0