53fb25efcf
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Check push changed paths (push) Successful in 19s
CI/CD Pipeline / Build Staging API Image (push) Successful in 41s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 48s
CI/CD Pipeline / Integration Tests (push) Successful in 3m10s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 3m17s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m30s
CI/CD Pipeline / Validate - Style (push) Successful in 4m17s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 59s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m33s
CI/CD Pipeline / Validate - Security (push) Successful in 7m12s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m38s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m13s
CI/CD Pipeline / Unit Tests (push) Successful in 10m11s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Staging E2E Tests (push) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (push) Failing after 26h14m3s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 26h24m21s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 26h19m47s
CI/CD Pipeline / PR Build Web Image (push) Failing after 26h23m44s
CI/CD Pipeline / PR Build API Image (push) Failing after 26h23m44s
CI/CD Pipeline / Deploy Production (push) Failing after 26h13m23s
CI/CD Pipeline / Build Production Web Image (push) Failing after 26h13m26s
CI/CD Pipeline / CI Gate (push) Failing after 26h13m25s
CI/CD Pipeline / Build Production API Image (push) Failing after 26h13m26s
CI/CD Pipeline / Canary Release to Production (push) Failing after 26h13m23s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Failing after 26h19m46s
CI/CD Pipeline / Frontend Lint (push) Failing after 26h23m37s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 26h23m45s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 26h19m46s
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
192 lines
6.5 KiB
Python
Executable File
192 lines
6.5 KiB
Python
Executable File
"""Asset InMemory Repository 实现"""
|
|
|
|
from datetime import UTC
|
|
|
|
from packages.domain import Asset
|
|
|
|
|
|
class InMemoryAssetRepository:
|
|
"""素材 InMemory 仓储实现(用于同步用例和测试)"""
|
|
|
|
def __init__(self):
|
|
self._assets: dict[str, Asset] = {}
|
|
|
|
def create(self, asset: Asset) -> Asset:
|
|
self._assets[asset.id] = asset
|
|
return asset
|
|
|
|
def get(self, asset_id: str) -> Asset | None:
|
|
return self._assets.get(asset_id)
|
|
|
|
def list_by_project(self, project_id: str) -> list[Asset]:
|
|
return [asset for asset in self._assets.values() if asset.project_id == project_id]
|
|
|
|
def list_by_library(self, library_id: str) -> list[Asset]:
|
|
return [asset for asset in self._assets.values() if asset.library_id == library_id]
|
|
|
|
def find_by_library(self, library_id: str) -> list[Asset]:
|
|
"""Alias for list_by_library to match the port interface."""
|
|
return self.list_by_library(library_id)
|
|
|
|
def find_by_library_and_file_type(self, library_id: str, file_type: str) -> list[Asset]:
|
|
return [
|
|
asset
|
|
for asset in self._assets.values()
|
|
if asset.library_id == library_id and asset.mime_type and asset.mime_type.startswith(file_type)
|
|
]
|
|
|
|
def update(self, asset: Asset) -> Asset:
|
|
self._assets[asset.id] = asset
|
|
return asset
|
|
|
|
def delete(self, asset_id: str) -> bool:
|
|
if asset_id in self._assets:
|
|
del self._assets[asset_id]
|
|
return True
|
|
return False
|
|
|
|
def batch_delete(self, asset_ids: list[str]) -> int:
|
|
"""批量删除素材(软删除,标记 status=deleted),返回实际影响数量。"""
|
|
from datetime import datetime
|
|
|
|
from packages.domain import AssetStatus
|
|
|
|
count = 0
|
|
for aid in asset_ids:
|
|
asset = self._assets.get(aid)
|
|
if asset and asset.status != AssetStatus.DELETED:
|
|
asset.status = AssetStatus.DELETED
|
|
asset.updated_at = datetime.now(UTC)
|
|
count += 1
|
|
return count
|
|
|
|
def batch_update_metadata(self, asset_ids: list[str], metadata_patch: dict[str, object]) -> int:
|
|
"""批量更新素材 metadata(合并 patch),返回实际影响数量。"""
|
|
from datetime import datetime
|
|
|
|
count = 0
|
|
for aid in asset_ids:
|
|
asset = self._assets.get(aid)
|
|
if asset:
|
|
asset.metadata = {**asset.metadata, **metadata_patch}
|
|
asset.updated_at = datetime.now(UTC)
|
|
count += 1
|
|
return count
|
|
|
|
def batch_add_tags(self, asset_ids: list[str], tag_ids: list[str]) -> int:
|
|
"""批量给素材添加标签(合并去重),返回实际影响数量。"""
|
|
from datetime import datetime
|
|
|
|
count = 0
|
|
for aid in asset_ids:
|
|
asset = self._assets.get(aid)
|
|
if asset:
|
|
changed = False
|
|
for tid in tag_ids:
|
|
if tid not in asset.tag_ids:
|
|
asset.tag_ids.append(tid)
|
|
changed = True
|
|
if changed:
|
|
asset.updated_at = datetime.now(UTC)
|
|
count += 1
|
|
return count
|
|
|
|
def batch_replace_tags(self, asset_ids: list[str], tag_ids: list[str]) -> int:
|
|
"""批量替换素材标签(全量覆盖),返回实际影响数量。"""
|
|
from datetime import datetime
|
|
|
|
count = 0
|
|
for aid in asset_ids:
|
|
asset = self._assets.get(aid)
|
|
if asset:
|
|
asset.tag_ids = list(tag_ids)
|
|
asset.updated_at = datetime.now(UTC)
|
|
count += 1
|
|
return count
|
|
|
|
def find_by_project(
|
|
self,
|
|
project_id: str,
|
|
skip: int = 0,
|
|
limit: int = 100,
|
|
) -> list[Asset]:
|
|
items = [a for a in self._assets.values() if a.project_id == project_id]
|
|
return items[skip : skip + limit]
|
|
|
|
def find_by_id(self, asset_id: str) -> Asset | None:
|
|
return self._assets.get(asset_id)
|
|
|
|
def find_by_tag_ids(
|
|
self,
|
|
tag_ids: list[str],
|
|
skip: int = 0,
|
|
limit: int = 100,
|
|
) -> list[Asset]:
|
|
"""查找包含所有指定标签的素材。"""
|
|
if not tag_ids:
|
|
return []
|
|
tag_set = set(tag_ids)
|
|
items = [a for a in self._assets.values() if tag_set.issubset(set(a.tag_ids))]
|
|
return items[skip : skip + limit]
|
|
|
|
def find_by_storage_key(self, storage_key: str) -> Asset | None:
|
|
"""按 storage_key 查找素材。"""
|
|
for asset in self._assets.values():
|
|
if asset.storage_key == storage_key:
|
|
return asset
|
|
return None
|
|
|
|
def find_by_library_and_file_hash(
|
|
self,
|
|
library_id: str,
|
|
file_hash: str,
|
|
) -> Asset | None:
|
|
"""按素材库 + 文件哈希查找已有素材(去重检测)。"""
|
|
if not file_hash:
|
|
return None
|
|
for asset in self._assets.values():
|
|
if asset.library_id == library_id and asset.file_hash == file_hash:
|
|
return asset
|
|
return None
|
|
|
|
def find_by_library_and_client_upload_id(
|
|
self,
|
|
library_id: str,
|
|
client_upload_id: str,
|
|
) -> Asset | None:
|
|
"""按素材库 + 客户端幂等 token 查找已有素材。"""
|
|
if not client_upload_id:
|
|
return None
|
|
for asset in self._assets.values():
|
|
if asset.library_id == library_id and getattr(asset, "client_upload_id", "") == client_upload_id:
|
|
return asset
|
|
return 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:
|
|
"""兜底去重:同库 + 同文件名(+同大小)且近期活动状态的素材。"""
|
|
from datetime import datetime, timedelta
|
|
|
|
if not name:
|
|
return None
|
|
from packages.domain import AssetStatus
|
|
|
|
cutoff = datetime.now(UTC) - timedelta(minutes=within_minutes)
|
|
candidates = [
|
|
a
|
|
for a in self._assets.values()
|
|
if a.library_id == library_id
|
|
and a.name == name
|
|
and a.status in (AssetStatus.UPLOADING, AssetStatus.PROCESSING)
|
|
and a.created_at >= cutoff
|
|
and (not file_size or file_size <= 0 or a.file_size == file_size)
|
|
]
|
|
if not candidates:
|
|
return None
|
|
return max(candidates, key=lambda a: a.created_at)
|