Files
xiaoxia-saas/tests/unit/test_ingest_validation.py
T
xiaoxia d25c7e4a75
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI Build & Deploy Pipeline / Build Staging Web Image (push) Successful in 26s
CI Build & Deploy Pipeline / Build Production API Image (push) Has been skipped
CI Build & Deploy Pipeline / Build Production Web Image (push) Has been skipped
CI Build & Deploy Pipeline / Build Production Worker 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 59s
CI/CD Pipeline / Unit Tests (push) Successful in 3m38s
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 3m52s
CI/CD Pipeline / Integration Tests (push) Successful in 1m28s
CI Build & Deploy Pipeline / Build Staging API Image (push) Successful in 6m3s
CI Build & Deploy Pipeline / Build Staging Worker Image (push) Successful in 7m40s
CI Build & Deploy Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 47s
CI Build & Deploy Pipeline / Staging API Integration Tests (push) Successful in 3m48s
CI Build & Deploy Pipeline / Staging E2E Tests (push) Failing after 4m51s
fix(worker): ingest先下载OSS文件再提取元数据,修复有效性校验全部误判ERROR (#489)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-07-18 00:25:13 +08:00

383 lines
14 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
# ── 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()