Files
xiaoxia-saas/tests/unit/test_ingest_validation.py
T
xiaoxia 854522986e
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI Build & Deploy Pipeline / Build Staging Web Image (push) Successful in 24s
CI Build & Deploy Pipeline / Build Production API Image (push) Has been skipped
CI Build & Deploy Pipeline / Build Production Worker Image (push) Has been skipped
CI Build & Deploy Pipeline / Build Production Web Image (push) Has been skipped
CI Build & Deploy Pipeline / Deploy Production (push) Has been skipped
CI Build & Deploy Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Frontend Lint (push) Successful in 1m1s
CI/CD Pipeline / Unit Tests (push) Successful in 3m12s
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 3m24s
CI/CD Pipeline / Integration Tests (push) Successful in 1m21s
CI Build & Deploy Pipeline / Build Staging API Image (push) Successful in 5m12s
CI Build & Deploy Pipeline / Build Staging Worker Image (push) Successful in 7m30s
CI Build & Deploy Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 3m2s
CI Build & Deploy Pipeline / Staging E2E Tests (push) Failing after 1m3s
CI Build & Deploy Pipeline / Staging API Integration Tests (push) Successful in 2m52s
fix: 自动选素材过滤HEVC等不支持编码 + ingest编码格式校验 (#494)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-07-18 11:27:16 +08:00

463 lines
16 KiB
Python
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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
def test_video_hevc_codec_invalid(self):
"""HEVC/H.265 编码视频 → 无效(渲染不支持)."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "hevc"},
"video",
)
assert result is False
def test_video_h265_codec_invalid(self):
"""h265 编码 → 无效."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "h265"},
"video",
)
assert result is False
def test_video_vp9_codec_invalid(self):
"""VP9 编码 → 无效."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "vp9"},
"video",
)
assert result is False
def test_video_av1_codec_invalid(self):
"""AV1 编码 → 无效."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "av1"},
"video",
)
assert result is False
def test_video_h264_codec_valid(self):
"""H.264 编码 → 有效."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "h264"},
"video",
)
assert result is True
def test_video_avc1_codec_valid(self):
"""avc1 编码(H.264变体) → 有效."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "avc1"},
"video",
)
assert result is True
def test_video_missing_codec_valid(self):
"""codec 字段缺失(存量数据) → 暂时放过,兼容历史."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0},
"video",
)
assert result is True
def test_video_codec_case_insensitive(self):
"""编码格式校验不区分大小写."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "HEVC"},
"video",
)
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 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")
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.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")
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()
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()