"""Issue #1714:POST /upload/direct/complete 幂等 + multipart 幂等。 覆盖: - 同 client_upload_id 重复 complete → 只建一条 asset、不重复派 ingest job - 同 file_hash 重复 complete → 返回已存在记录 - 旧客户端不传 hash/token:近期同库同名 processing 占位 → 兜底幂等返回 - 旧客户端不传 hash/token:READY 历史同名 → 不兜底(正常新建) - 兜底窗口外(>30 分钟)→ 不兜底 - 旧仓储(无新方法)鸭子类型降级 → 不报错、正常新建 - 重复 complete 时即使 OSS 已无文件(file_exists=False)也返回已存在记录 (模拟 complete 超时后 OSS 侧对象已过期/清理,重试仍不重复建库) - multipart 上传:同 client_upload_id 重复提交 → 第二次直接 duplicated,不再传 OSS """ from __future__ import annotations import os import sys from datetime import datetime, timedelta, timezone from pathlib import Path from unittest.mock import MagicMock os.environ.setdefault("JWT_SECRET_KEY", "unit-test-secret-key-for-testing") os.environ.setdefault("DATABASE_URL", "sqlite:///test.db") sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "apps" / "api")) from fastapi import FastAPI # noqa: E402 from fastapi.testclient import TestClient # noqa: E402 from packages.domain import Asset, AssetLibrary, AssetLibraryKind, AssetStatus, IngestJob, Project # noqa: E402 class StubProjectRepository: def __init__(self, projects: dict | None = None): self._projects = projects or {} def get(self, project_id: str): return self._projects.get(project_id) def find_by_id(self, project_id: str): return self._projects.get(project_id) class StubAssetLibraryRepository: def __init__(self, libraries: dict | None = None): self._libraries = libraries or {} def find_by_project(self, project_id: str, kind=None) -> list: return list(self._libraries.values()) class StubAssetRepository: """支持三种幂等查询的内存仓储,并统计 create 次数。""" def __init__(self, assets: list[Asset] | None = None): self._assets = list(assets or []) self.created: list[Asset] = [] def find_by_library_and_file_hash(self, library_id: str, file_hash: str) -> Asset | None: if not file_hash: return None return next((a for a in self._assets if a.library_id == library_id and a.file_hash == file_hash), None) def find_by_library_and_client_upload_id(self, library_id: str, client_upload_id: str) -> Asset | None: if not client_upload_id: return None return next( (a for a in self._assets if a.library_id == library_id and a.client_upload_id == client_upload_id), None, ) def find_recent_active_by_library_and_name( self, library_id: str, name: str, within_minutes: int = 30, file_size: int = 0 ) -> Asset | None: # 严格模式(#1714):大小未知(0)直接不命中,宁可漏判不可误杀 if not file_size or file_size <= 0: return None cutoff = datetime.now(timezone.utc) - timedelta(minutes=within_minutes) candidates = [ a for a in self._assets if a.library_id == library_id and a.name == name and a.status in (AssetStatus.UPLOADING, AssetStatus.PROCESSING) and a.created_at >= cutoff and a.file_size == file_size ] return max(candidates, key=lambda a: a.created_at) if candidates else None def create(self, asset: Asset) -> Asset: self._assets.append(asset) self.created.append(asset) return asset def update(self, asset: Asset) -> Asset: return asset class LegacyStubAssetRepository: """旧仓储:只有 file_hash 去重,没有新方法(鸭子类型降级验证)。""" def __init__(self, assets: list[Asset] | None = None): self._assets = list(assets or []) self.created: list[Asset] = [] def find_by_library_and_file_hash(self, library_id: str, file_hash: str) -> Asset | None: if not file_hash: return None return next((a for a in self._assets if a.library_id == library_id and a.file_hash == file_hash), None) def create(self, asset: Asset) -> Asset: self._assets.append(asset) self.created.append(asset) return asset class StubIngestJobRepository: def __init__(self): self._jobs: dict[str, IngestJob] = {} self.created_count = 0 def create(self, job: IngestJob) -> IngestJob: self._jobs[job.id] = job self.created_count += 1 return job def get(self, job_id: str) -> IngestJob | None: return self._jobs.get(job_id) def update(self, job: IngestJob) -> IngestJob: self._jobs[job.id] = job return job def _make_project() -> Project: return Project(id="proj-1", name="Test Project", owner_user_id="user-1") def _make_library() -> AssetLibrary: return AssetLibrary(id="lib-1", name="Test Library", project_id="proj-1", kind=AssetLibraryKind.VIDEO) def _build_app(asset_repo=None, ingest_repo=None, storage=None): from app.api.routes.upload import router from app.auth import AuthenticatedUser, get_current_user from app.core.storage import get_storage_service from app.dependencies import ( get_asset_library_repository, get_asset_repository, get_ingest_job_repository, get_project_repository, ) app = FastAPI() app.include_router(router, prefix="/api/v1") project_repo = StubProjectRepository({"proj-1": _make_project()}) library_repo = StubAssetLibraryRepository({"lib-1": _make_library()}) asset_repo = asset_repo or StubAssetRepository() ingest_repo = ingest_repo or StubIngestJobRepository() storage = storage or MagicMock() storage.is_configured = True storage._normalize_storage_key = lambda key: key storage.file_exists = MagicMock(return_value=True) storage.upload_file = MagicMock(return_value="https://oss.example.com/file.mp4") storage.get_url = MagicMock(return_value="https://oss.example.com/file.mp4") mock_user = MagicMock(spec=AuthenticatedUser) mock_user.id = "user-1" mock_user.user = MagicMock(id="user-1") mock_user.email = "test@example.com" app.dependency_overrides[get_current_user] = lambda: mock_user app.dependency_overrides[get_project_repository] = lambda: project_repo app.dependency_overrides[get_asset_library_repository] = lambda: library_repo app.dependency_overrides[get_asset_repository] = lambda: asset_repo app.dependency_overrides[get_ingest_job_repository] = lambda: ingest_repo app.dependency_overrides[get_storage_service] = lambda: storage return app, asset_repo, ingest_repo, storage def _client(**kwargs): app, asset_repo, ingest_repo, storage = _build_app(**kwargs) return TestClient(app), asset_repo, ingest_repo, storage COMPLETE_BODY = { "project_id": "proj-1", "library_id": "lib-1", "storage_key": "uploads/abc/IMG_2282.MOV", } class TestDirectCompleteIdempotency: def test_same_client_upload_id_creates_single_asset_and_job(self): """同一 client_upload_id 连发两次 complete:只建 1 条 asset、1 个 job。""" client, asset_repo, ingest_repo, _ = _client() body = {**COMPLETE_BODY, "client_upload_id": "up-token-1", "file_size": 12345} r1 = client.post("/api/v1/direct/complete", json=body) r2 = client.post("/api/v1/direct/complete", json={**body, "storage_key": "uploads/zzz/IMG_2282.MOV"}) assert r1.status_code == 200 and r2.status_code == 200 b1, b2 = r1.json(), r2.json() assert b1["duplicated"] is False assert b2["duplicated"] is True assert b1["asset_id"] == b2["asset_id"] assert len(asset_repo.created) == 1 assert ingest_repo.created_count == 1 # 第二次返回的是已存在记录(其 storage_key 为第一次的 key) assert b2["storage_key"] == "uploads/abc/IMG_2282.MOV" def test_same_file_hash_returns_existing(self): """同 file_hash(不同 token)重复 complete → 返回已存在记录。""" client, asset_repo, ingest_repo, _ = _client() body1 = {**COMPLETE_BODY, "file_hash": "h" * 32, "client_upload_id": "tok-a"} body2 = { **COMPLETE_BODY, "storage_key": "uploads/def/IMG_2282.MOV", "file_hash": "h" * 32, "client_upload_id": "tok-b", } client.post("/api/v1/direct/complete", json=body1) r2 = client.post("/api/v1/direct/complete", json=body2) assert r2.json()["duplicated"] is True assert len(asset_repo.created) == 1 assert ingest_repo.created_count == 1 def test_fallback_dedup_when_no_hash_no_token(self): """旧客户端不传 hash/token:近期同库同名 processing 占位 → 兜底幂等。 模拟 complete 超时重试:第一次已建好占位,第二次(OSS 重传拿到新 key) 不应再建第二条。 """ client, asset_repo, ingest_repo, _ = _client() # 第一次 complete(旧客户端无 token/hash,但 file_size 可知) r1 = client.post( "/api/v1/direct/complete", json={**COMPLETE_BODY, "file_size": 5_000_000}, ) assert r1.json()["duplicated"] is False # 重试:重新 prepare 产生新 storage_key(仅 uuid 目录不同,文件名一致—— # 前端重试传的是同一个 File),且近期;同大小才允许兜底命中 r2 = client.post( "/api/v1/direct/complete", json={ **COMPLETE_BODY, "storage_key": "uploads/retry/IMG_2282.MOV", "file_size": 5_000_000, }, ) assert r2.status_code == 200 assert r2.json()["duplicated"] is True assert r2.json()["asset_id"] == r1.json()["asset_id"] assert len(asset_repo.created) == 1 assert ingest_repo.created_count == 1 def test_fallback_dedup_skipped_when_file_size_unknown(self): """file_size=0(未知)时不允许仅凭同名 + processing 判重,直接放行(#1714)。 根因场景:complete 没传 file_size,30 分钟内同名占位(如 iPhone 的 IMG_2285.MOV)会把内容/大小全新的视频误判为重复跳过。 """ client, asset_repo, _ingest_repo, _ = _client() r1 = client.post( "/api/v1/direct/complete", json={**COMPLETE_BODY, "file_size": 0}, ) assert r1.json()["duplicated"] is False # 第二个全新视频:同名(IMG_2285.MOV)、无 hash/token、file_size 仍未知 r2 = client.post( "/api/v1/direct/complete", json={**COMPLETE_BODY, "storage_key": "uploads/retry2/IMG_2282.MOV", "file_size": 0}, ) assert r2.status_code == 200 assert r2.json()["duplicated"] is False # 不能误杀 assert len(asset_repo.created) == 2 # 两条记录,放行新上传 def test_fallback_dedup_skipped_when_same_name_but_different_size(self): """同名但 file_size 不同 → 不判重,正常建记录(#1714)。""" client, asset_repo, _ingest_repo, _ = _client() r1 = client.post( "/api/v1/direct/complete", json={**COMPLETE_BODY, "file_size": 5_000_000}, ) assert r1.json()["duplicated"] is False r2 = client.post( "/api/v1/direct/complete", json={ **COMPLETE_BODY, "storage_key": "uploads/retry3/IMG_2282.MOV", "file_size": 9_999_999, # 同名但大小完全不同的新视频 }, ) assert r2.status_code == 200 assert r2.json()["duplicated"] is False assert len(asset_repo.created) == 2 def test_fallback_dedup_skipped_when_hash_present_even_if_name_size_match(self): """file_hash 非空且 hash 未命中时,不允许退回同名兜底(#1714)。 hash 已能代表内容:同名同大小但 hash 不同是真实的新内容,必须放行。 """ client, asset_repo, _ingest_repo, _ = _client() # 第一次:某 hash 的视频 r1 = client.post( "/api/v1/direct/complete", json={ **COMPLETE_BODY, "file_hash": "a" * 64, "client_upload_id": "tok-1", "file_size": 5_000_000, }, ) assert r1.json()["duplicated"] is False # 第二次:同名同大小但 hash 不同(新视频内容不同); # 注意 client_upload_id 也必须不同,否则会先被 token 命中 r2 = client.post( "/api/v1/direct/complete", json={ **COMPLETE_BODY, "storage_key": "uploads/retry4/IMG_2282.MOV", "file_hash": "b" * 64, "client_upload_id": "tok-2", "file_size": 5_000_000, }, ) assert r2.status_code == 200 assert r2.json()["duplicated"] is False assert len(asset_repo.created) == 2 def test_fallback_dedup_ignores_ready_history(self): """READY 历史同名素材不触发兜底(允许用户再次上传同名文件)。""" ready = Asset( id="ready-1", project_id="proj-1", library_id="lib-1", name="IMG_2282.MOV", storage_key="uploads/old/IMG_2282.MOV", mime_type="video/quicktime", status=AssetStatus.READY, ) client, asset_repo, ingest_repo, _ = _client(asset_repo=StubAssetRepository([ready])) r = client.post("/api/v1/direct/complete", json=COMPLETE_BODY) assert r.status_code == 200 assert r.json()["duplicated"] is False assert len(asset_repo.created) == 1 def test_fallback_dedup_window_expired(self): """占位记录超过 30 分钟 → 不再兜底(视为孤儿,正常新建)。""" stale = Asset( id="stale-1", project_id="proj-1", library_id="lib-1", name="IMG_2282.MOV", storage_key="uploads/stale/IMG_2282.MOV", mime_type="video/quicktime", status=AssetStatus.PROCESSING, ) stale.created_at = datetime.now(timezone.utc) - timedelta(minutes=45) client, asset_repo, ingest_repo, _ = _client(asset_repo=StubAssetRepository([stale])) r = client.post("/api/v1/direct/complete", json=COMPLETE_BODY) assert r.status_code == 200 assert r.json()["duplicated"] is False assert len(asset_repo.created) == 1 def test_legacy_repo_without_new_methods_still_works(self): """旧仓储没有新幂等方法 → 鸭子类型降级,不报错、正常创建。""" client, asset_repo, ingest_repo, _ = _client(asset_repo=LegacyStubAssetRepository()) r = client.post( "/api/v1/direct/complete", json={**COMPLETE_BODY, "client_upload_id": "tok-x", "file_hash": "f" * 32}, ) assert r.status_code == 200 assert r.json()["duplicated"] is False assert len(asset_repo.created) == 1 def test_duplicate_complete_returns_existing_even_if_oss_missing(self): """重复 complete 幂等检查先于 OSS file_exists: 第一次成功建占位后,重试时即使 OSS 对象已不存在(file_exists=False), 也必须返回已存在记录而不是 404/重复建库。""" client, _, _, storage = _client() body = {**COMPLETE_BODY, "client_upload_id": "tok-oss-gone"} r1 = client.post("/api/v1/direct/complete", json=body) assert r1.status_code == 200 storage.file_exists = MagicMock(return_value=False) r2 = client.post( "/api/v1/direct/complete", json={**body, "storage_key": "uploads/retry2/IMG_2282.MOV"}, ) assert r2.status_code == 200 assert r2.json()["duplicated"] is True assert r2.json()["asset_id"] == r1.json()["asset_id"] class TestMultipartUploadIdempotency: def test_same_client_upload_id_second_submit_deduplicated(self): """multipart 重复提交同 token:第二次直接 duplicated,不再上传 OSS。""" client, asset_repo, ingest_repo, storage = _client() def _post(): return client.post( "/api/v1", data={"project_id": "proj-1", "library_id": "lib-1", "client_upload_id": "mp-tok-1"}, files={"file": ("IMG_2282.MOV", b"fake-mov-data", "video/quicktime")}, ) r1 = _post() r2 = _post() assert r1.json()["duplicated"] is False assert r2.json()["duplicated"] is True assert r2.json()["asset_id"] == r1.json()["asset_id"] assert len(asset_repo.created) == 1 assert ingest_repo.created_count == 1 # OSS 上传只发生一次(第二次在幂等检查处直接返回) assert storage.upload_file.call_count == 1