fix(worker): ingest增加素材有效性校验,坏文件标记为ERROR #488

Merged
auto-approve-bot merged 1 commits from fix/ingest-validation into develop 2026-07-17 22:20:57 +08:00
2 changed files with 400 additions and 20 deletions
+133 -20
View File
@@ -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
+267
View File
@@ -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=ERRORjob=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=READYjob=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()