From 355e25b2aa23fd88300eee566be957cd39b41f5a Mon Sep 17 00:00:00 2001 From: CI Bot Date: Fri, 17 Jul 2026 22:15:22 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20ingest=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E7=B4=A0=E6=9D=90=E6=9C=89=E6=95=88=E6=80=A7=E6=A0=A1=E9=AA=8C?= =?UTF-8?q?=EF=BC=8C=E5=9D=8F=E6=96=87=E4=BB=B6=E6=A0=87=E8=AE=B0=E4=B8=BA?= =?UTF-8?q?ERROR=E4=B8=8D=E5=86=8D=E4=BB=A5READY=E5=85=A5=E5=BA=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题:ingest流程中ffprobe/Pillow失败时仍然创建status=READY的asset, 导致27字节文本文件也能伪装成有效视频进入素材库, 一键生成自动选到坏素材直接渲染失败。 修复: - extract_media_metadata返回(metadata, success)二元组 - 新增_is_valid_media校验函数: * 视频:ffprobe成功 + size>=1KB + duration>0 * 音频:ffprobe成功 + size>=100B + duration>0 * 图片:Pillow验证通过 + size>=100B + width>0 + height>0 - 校验失败的asset标记为ERROR状态,ingest job标记为FAILED - 新增14个单元测试:有效性判断+返回值签名+完整ingest流程 - 补充音频元数据提取(之前只有视频和图片) --- apps/worker/worker_app/tasks/ingest.py | 153 ++++++++++++-- tests/unit/test_ingest_validation.py | 267 +++++++++++++++++++++++++ 2 files changed, 400 insertions(+), 20 deletions(-) create mode 100755 tests/unit/test_ingest_validation.py diff --git a/apps/worker/worker_app/tasks/ingest.py b/apps/worker/worker_app/tasks/ingest.py index 392c2c70b..47879b913 100755 --- a/apps/worker/worker_app/tasks/ingest.py +++ b/apps/worker/worker_app/tasks/ingest.py @@ -15,6 +15,12 @@ from packages.domain import Asset, AssetStatus, IngestJobStatus logger = get_task_logger(__name__) +# 最小有效文件大小(字节):小于此值的直接判为无效,避免文本/空文件伪装成媒体 +MIN_VIDEO_FILE_SIZE = 1024 # 1KB +MIN_AUDIO_FILE_SIZE = 100 # 100B +MIN_IMAGE_FILE_SIZE = 100 # 100B + + def _safe_parse_fps(fps_str: str) -> float: """Safely parse fps from a fraction string like \"30/1\" or \"30000/1001\".""" try: @@ -29,7 +35,7 @@ def _safe_parse_fps(fps_str: str) -> float: return 0.0 -def extract_media_metadata(file_url: str, media_type: str) -> dict: +def extract_media_metadata(file_url: str, media_type: str) -> tuple[dict, bool]: """ 提取媒体文件的元数据。 @@ -38,9 +44,12 @@ def extract_media_metadata(file_url: str, media_type: str) -> dict: media_type: 媒体类型 (video, audio, image) Returns: - 提取的元数据字典,失败时返回空字典 + (metadata_dict, success) + - metadata_dict: 提取的元数据字典,失败时返回空字典 + - success: 是否成功提取到有效元数据 """ metadata = {} + success = False try: if media_type == "video": @@ -80,23 +89,66 @@ def extract_media_metadata(file_url: str, media_type: str) -> dict: metadata["size_bytes"] = int(format_info.get("size", 0)) metadata["bitrate"] = int(format_info.get("bit_rate", 0)) + # 有效性判断:有视频流 + duration > 0 + size > 0 + has_video_stream = any(s.get("codec_type") == "video" for s in probe_data.get("streams", [])) + has_audio_stream = any(s.get("codec_type") == "audio" for s in probe_data.get("streams", [])) + if (has_video_stream or has_audio_stream) and metadata.get("duration", 0) > 0: + success = True + except Exception as e: logger.warning("视频元数据提取失败: %s", e) + elif media_type == "audio": + from video_processing.ffmpeg_utils import run_ffprobe + + cmd = [ + "ffprobe", + "-v", + "quiet", + "-print_format", + "json", + "-show_format", + "-show_streams", + file_url, + ] + try: + stdout, _ = run_ffprobe(cmd, timeout=30) + import json as json_lib + + probe_data = json_lib.loads(stdout) + + has_audio_stream = any(s.get("codec_type") == "audio" for s in probe_data.get("streams", [])) + format_info = probe_data.get("format", {}) + metadata["duration"] = float(format_info.get("duration", 0)) # type: ignore[assignment] + metadata["size_bytes"] = int(format_info.get("size", 0)) + metadata["bitrate"] = int(format_info.get("bit_rate", 0)) + metadata["codec"] = next( + (s.get("codec_name", "") for s in probe_data.get("streams", []) if s.get("codec_type") == "audio"), + "", + ) + + if has_audio_stream and metadata.get("duration", 0) > 0: + success = True + + except Exception as e: + logger.warning("音频元数据提取失败: %s", e) + elif media_type == "image": # 使用 Pillow 提取图片元数据 try: from PIL import Image with Image.open(file_url) as img: - metadata["width"] = img.width - metadata["height"] = img.height - metadata["format"] = img.format - metadata["mode"] = img.mode - if hasattr(img, "_getexif") and img._getexif(): - exif = img._getexif() - if exif: - metadata["exif"] = {k: str(v) for k, v in exif.items() if isinstance(v, (str, int, float))} # type: ignore[assignment] + img.verify() # 验证文件完整性 + # verify后需要重新打开才能读尺寸 + with Image.open(file_url) as img2: + metadata["width"] = img2.width + metadata["height"] = img2.height + metadata["format"] = img2.format + metadata["mode"] = img2.mode + if img2.width > 0 and img2.height > 0: + success = True + except ImportError: logger.warning("Pillow not available for image metadata extraction") except Exception as e: @@ -109,7 +161,32 @@ def extract_media_metadata(file_url: str, media_type: str) -> dict: except Exception as e: logger.warning(f"Failed to extract metadata: {e}") - return metadata + return metadata, success + + +def _is_valid_media(metadata: dict, media_type: str) -> bool: + """根据元数据判断文件是否为有效媒体文件。 + + Args: + metadata: extract_media_metadata 返回的元数据 + media_type: 媒体类型 + + Returns: + True 表示文件有效 + """ + size = int(metadata.get("size_bytes", 0)) + + if media_type == "video": + duration = float(metadata.get("duration", 0)) + return size >= MIN_VIDEO_FILE_SIZE and duration > 0 + if media_type == "audio": + duration = float(metadata.get("duration", 0)) + return size >= MIN_AUDIO_FILE_SIZE and duration > 0 + if media_type == "image": + width = int(metadata.get("width", 0)) + height = int(metadata.get("height", 0)) + return size >= MIN_IMAGE_FILE_SIZE and width > 0 and height > 0 + return False @celery_app.task(name="worker.ingest_asset") @@ -150,17 +227,53 @@ def ingest_asset(job_id: str) -> dict: elif mime_type.startswith("audio/"): media_type = "audio" - # Extract metadata (returns empty dict on failure) + # Extract metadata storage_url = job.storage_key # Assuming storage_key is usable as URL/path - metadata = extract_media_metadata(storage_url, media_type) + metadata, extract_success = extract_media_metadata(storage_url, media_type) - # Fill in defaults if metadata extraction failed - if not metadata: - metadata = { - "duration": 0, - "width": 0, - "height": 0, - "size_bytes": 0, + # 有效性校验:ffprobe/Pillow 必须成功,且文件大小/时长/尺寸满足最小要求 + is_valid = extract_success and _is_valid_media(metadata, media_type) + + if not is_valid: + # 文件无效,创建 ERROR 状态的 asset 并标记 job 失败 + error_reason = "metadata extraction failed" if not extract_success else "media validation failed" + logger.warning( + "素材有效性校验失败,标记为ERROR: job_id=%s storage_key=%s media_type=%s reason=%s", + job_id, + job.storage_key, + media_type, + error_reason, + ) + + asset = Asset.create( + project_id=job.project_id, + library_id=job.library_id, + name=filename, + storage_key=job.storage_key, + mime_type=mime_type, + metadata={"ingest_error": error_reason}, + file_size=int(metadata.get("size_bytes", 0)), + duration=float(metadata.get("duration", 0)), + width=int(metadata.get("width", 0)), + height=int(metadata.get("height", 0)), + status=AssetStatus.ERROR, + file_hash=job.file_hash, + ) + asset_repo.create(asset) + + # Update job status to FAILED + job.status = IngestJobStatus.FAILED + job.error_message = f"Invalid media file: {error_reason}" + job.result_asset_id = asset.id + job.updated_at = datetime.now(timezone.utc) + job_repo.update(job) + db.commit() + + return { + "status": "failed", + "job_id": job.id, + "asset_id": asset.id, + "error": error_reason, } # Create Asset diff --git a/tests/unit/test_ingest_validation.py b/tests/unit/test_ingest_validation.py new file mode 100755 index 000000000..31877e6b0 --- /dev/null +++ b/tests/unit/test_ingest_validation.py @@ -0,0 +1,267 @@ +"""ingest 素材有效性校验单元测试 — 坏文件应该标记为ERROR,不能以READY入库.""" + +from __future__ import annotations + +import os +import sys +from pathlib import Path +from unittest.mock import MagicMock, patch + +# 在 import worker_app 模块前 mock 掉数据库连接和celery +_mock_db_module = MagicMock() +_mock_db_module.SessionLocal = MagicMock() +sys.modules["worker_app.db"] = _mock_db_module +sys.modules["worker_app.core.config"] = MagicMock() + +# mock celery_app.task使其返回原函数(装饰器透传) +_mock_celery_module = MagicMock() + + +def _passthrough_decorator(*args, **kwargs): + if len(args) == 1 and callable(args[0]): + return args[0] + return lambda f: f + + +_mock_celery_module.celery_app.task = MagicMock(side_effect=_passthrough_decorator) +sys.modules["worker_app.celery_app"] = _mock_celery_module + +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "apps" / "worker")) + +import pytest + +# ── _is_valid_media 纯函数测试 ────────────────────────────────────────────── + + +class TestIsValidMedia: + """_is_valid_media 有效性判断逻辑.""" + + def test_video_valid_normal(self): + """正常视频:有大小有时长 → 有效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 5 * 1024 * 1024, "duration": 30.5}, "video") + assert result is True + + def test_video_zero_size_invalid(self): + """视频0字节 → 无效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 0, "duration": 30}, "video") + assert result is False + + def test_video_too_small_invalid(self): + """视频小于1KB → 无效(文本文件伪装).""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 27, "duration": 0}, "video") + assert result is False + + def test_video_zero_duration_invalid(self): + """视频时长为0 → 无效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 1024 * 1024, "duration": 0}, "video") + assert result is False + + def test_audio_valid_normal(self): + """正常音频 → 有效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 3 * 1024 * 1024, "duration": 180.0}, "audio") + assert result is True + + def test_audio_zero_size_invalid(self): + """音频0字节 → 无效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 0, "duration": 60}, "audio") + assert result is False + + def test_image_valid_normal(self): + """正常图片 → 有效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 200 * 1024, "width": 1920, "height": 1080}, "image") + assert result is True + + def test_image_zero_dimension_invalid(self): + """图片0尺寸 → 无效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 1024, "width": 0, "height": 0}, "image") + assert result is False + + def test_unknown_media_type_invalid(self): + """未知媒体类型 → 无效.""" + from worker_app.tasks.ingest import _is_valid_media + + result = _is_valid_media({"size_bytes": 1000}, "text") + assert result is False + + +# ── extract_media_metadata 返回值签名测试 ──────────────────────────────────── + + +class TestExtractMediaMetadataReturnValue: + """extract_media_metadata 返回 (dict, bool) 元组.""" + + def test_returns_tuple_with_two_elements(self): + """返回值是 (metadata, success) 二元组.""" + from worker_app.tasks.ingest import extract_media_metadata + + # 用不存在的文件路径,ffprobe必然失败 + metadata, success = extract_media_metadata("/nonexistent/file.mp4", "video") + assert isinstance(metadata, dict) + assert isinstance(success, bool) + assert success is False + + def test_video_ffprobe_failure_returns_false(self): + """ffprobe解析失败 → success=False.""" + from worker_app.tasks.ingest import extract_media_metadata + + _, success = extract_media_metadata("/nonexistent/fake.mp4", "video") + assert success is False + + def test_image_pillow_failure_returns_false(self): + """Pillow打不开 → success=False.""" + from worker_app.tasks.ingest import extract_media_metadata + + _, success = extract_media_metadata("/nonexistent/fake.jpg", "image") + assert success is False + + +# ── 完整 ingest 流程测试(mock ffprobe + mock db) ──────────────────────── + + +class TestIngestAssetValidation: + """完整ingest流程的有效性校验:坏文件→ERROR状态,好文件→READY状态.""" + + def test_invalid_video_marked_as_error(self): + """ffprobe返回空数据(坏文件)→ asset.status=ERROR,job=FAILED.""" + # 内存数据库 + 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 + job = IngestJob( + id="job-invalid-001", + project_id="proj-1", + library_id="lib-1", + storage_key="uploads/test/bad.mp4", + status=IngestJobStatus.PENDING, + ) + 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 + + result = ingest_asset("job-invalid-001") + + assert result["status"] == "failed" + assert "asset_id" in result + + # 验证asset状态是ERROR + asset = asset_repo.get(result["asset_id"]) + assert asset is not None + assert asset.status == AssetStatus.ERROR + assert asset.file_size == 27 + assert asset.duration == 0 + + # 验证job状态是FAILED + updated_job = job_repo.get("job-invalid-001") + assert updated_job.status == IngestJobStatus.FAILED + assert "Invalid media file" in updated_job.error_message + + db.close() + + def test_valid_video_marked_as_ready(self): + """正常视频 → asset.status=READY,job=COMPLETED.""" + 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-valid-001", + project_id="proj-1", + library_id="lib-1", + storage_key="uploads/test/good.mp4", + status=IngestJobStatus.PENDING, + ) + 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 + + result = ingest_asset("job-valid-001") + + assert result["status"] == "completed" + assert "asset_id" in result + + asset = asset_repo.get(result["asset_id"]) + assert asset is not None + assert asset.status == AssetStatus.READY + assert asset.duration == 30.5 + assert asset.width == 1920 + assert asset.height == 1080 + + updated_job = job_repo.get("job-valid-001") + assert updated_job.status == IngestJobStatus.COMPLETED + + db.close() -- 2.54.0