fix(worker): ingest增加素材有效性校验,坏文件标记为ERROR #488
@@ -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
|
||||
|
||||
Executable
+267
@@ -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()
|
||||
Reference in New Issue
Block a user