"""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=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 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=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.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()