Files
xiaoxia-saas/tests/unit/test_ingest_validation.py
xiaoxia edb118a5a0
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Frontend Lint (push) Successful in 2m43s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 3m23s
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 3m27s
CI/CD Pipeline / Unit Tests (push) Successful in 4m3s
CI/CD Pipeline / Integration Tests (push) Successful in 1m55s
fix: 修复ingest测试MagicMock写入SQLite导致Unit Tests失败 (#601)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-07-19 17:40:44 +08:00

474 lines
17 KiB
Python
Executable File
Raw Permalink 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_valid(self):
"""HEVC/H.265 编码视频 → 有效(ingest层不拦截,渲染层统一转码)."""
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 True
def test_video_h265_codec_valid(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 True
def test_video_vp9_codec_valid(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 True
def test_video_av1_codec_valid(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 True
def test_video_unknown_codec_valid(self):
"""非白名单编码 → 仍有效(只打日志不拦截,渲染层负责统一转码)."""
from worker_app.tasks.ingest import _is_valid_media
result = _is_valid_media(
{"size_bytes": 10 * 1024 * 1024, "duration": 30.0, "codec": "mpeg4"},
"video",
)
assert result is True
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):
"""编码格式校验不区分大小写,HEVC大写也识别为有效."""
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 True
# ── 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):
with patch("video_processing.thumbnail_generator.generate_and_upload_thumbnail", return_value=None):
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()