Files
xiaoxia-saas/packages/adapters/sqlalchemy_impl/asset_repository.py
T
saas-backend-agent fec656941d
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 6s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 4s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 24s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 24s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m41s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 1m57s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m57s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m44s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 2m50s
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
AI Code Review / AI Code Review (pull_request) Has been cancelled
PR Automation / Auto Approve on CI Green (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 241h52m35s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 241h52m38s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 241h52m43s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 241h52m49s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 241h52m49s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 241h53m9s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 241h53m9s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 241h53m15s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 241h53m26s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 241h53m16s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 242h27m12s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 242h27m25s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 242h27m45s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 242h27m52s
fix(#1714): 同名兜底去重误杀新视频 + ingest 链路孤儿清理/启动恢复
根因(staging 实证 IMG_2285.MOV file_size=0 卡 processing 持续误杀):
direct complete 不传 file_size,find_recent_active_by_library_and_name
的大小校验 if file_size>0 不生效,30 分钟内同名视频(iPhone IMG_xxxx.MOV)
即使内容全新也被同名兜底误判重复跳过。

兜底去重收紧(宁可漏判不可误杀):
- _find_duplicate_asset:file_hash/client_upload_id 非空时不走同名兜底
  (hash 已代表内容);走到兜底必须 file_size>0 且与记录大小严格一致
- find_recent_active_by_library_and_name:file_size=0 直接返回 None,
  大小条件改为 SQL 内严格等值匹配
- complete 的 _create_pending_asset 补传 file_size(之前占位记录大小永远 0)

ingest 链路孤儿兜底(此前只有 generation 链路有清理):
- 新增 packages/application/ingest_orphan_cleanup.py:
  - processing>60min / pending>90min 的 ingest_job 标 failed,关联
    processing/uploading asset 联动标 error;无 job 关联超 120min 孤儿
    占位 asset 也标 error(beat 每 10 分钟巡检)
  - worker 启动恢复:processing 超 10 分钟的 job CAS 重置 pending 并
    重新派单(Redis SET NX 锁互斥双 worker,旧消息重投由执行前守卫丢弃)
- celery 配置:task_reject_on_worker_lost=True;
  broker visibility_timeout=4h(acks_late 下长转码任务不被误重投)

单测:同名不同大小放行 / hash 非空不走兜底 / file_size=0 放行 /
同名同大小 processing 才判重;孤儿清理 9 例、启动恢复 4 例、
仓储严格守卫 6 例、beat 串联 2 例。diff coverage 100%。

前端配合(前端工程师):completeDirectUpload 请求体补 file_size=file.size。
2026-09-06 13:29:59 +08:00

514 lines
20 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.
import json
from datetime import datetime, timezone
from sqlalchemy.orm import Session
from packages.adapters.sqlalchemy_impl.models import AssetModel, AssetTagModel
from packages.domain import Asset, AssetStatus, ClassificationStatus
class SQLAlchemyAssetRepository:
def __init__(self, session: Session):
self.session = session
def find_by_library(
self,
library_id: str,
skip: int = 0,
limit: int = 100,
status: list[str] | None = None,
) -> list[Asset]:
query = self.session.query(AssetModel).filter(AssetModel.asset_library_id == library_id)
if status:
query = query.filter(AssetModel.status.in_(status))
models = query.order_by(AssetModel.created_at.desc()).offset(skip).limit(limit).all()
return [self._to_domain(model) for model in models]
def find_by_project(
self,
project_id: str,
skip: int = 0,
limit: int = 100,
status: list[str] | None = None,
) -> list[Asset]:
query = self.session.query(AssetModel).filter(AssetModel.project_id == project_id)
if status:
query = query.filter(AssetModel.status.in_(status))
models = query.order_by(AssetModel.created_at.desc()).offset(skip).limit(limit).all()
return [self._to_domain(model) for model in models]
def find_by_library_and_file_type(
self,
library_id: str,
file_type: str,
skip: int = 0,
limit: int = 100,
status: list[str] | None = None,
) -> list[Asset]:
query = self.session.query(AssetModel).filter(
AssetModel.asset_library_id == library_id, AssetModel.file_type == file_type
)
if status:
query = query.filter(AssetModel.status.in_(status))
models = query.order_by(AssetModel.created_at.desc()).offset(skip).limit(limit).all()
return [self._to_domain(model) for model in models]
def count_by_library_and_file_type(
self,
library_id: str,
file_type: str,
status: list[str] | None = None,
) -> int:
query = self.session.query(AssetModel).filter(
AssetModel.asset_library_id == library_id, AssetModel.file_type == file_type
)
if status:
query = query.filter(AssetModel.status.in_(status))
return query.count()
def find_by_project_and_file_type(
self,
project_id: str,
file_type: str,
skip: int = 0,
limit: int = 100,
status: list[str] | None = None,
) -> list[Asset]:
query = self.session.query(AssetModel).filter(
AssetModel.project_id == project_id, AssetModel.file_type == file_type
)
if status:
query = query.filter(AssetModel.status.in_(status))
models = query.order_by(AssetModel.created_at.desc()).offset(skip).limit(limit).all()
return [self._to_domain(model) for model in models]
def count_by_project_and_file_type(
self,
project_id: str,
file_type: str,
status: list[str] | None = None,
) -> int:
query = self.session.query(AssetModel).filter(
AssetModel.project_id == project_id, AssetModel.file_type == file_type
)
if status:
query = query.filter(AssetModel.status.in_(status))
return query.count()
def find_by_id(self, asset_id: str) -> Asset | None:
model = self.session.query(AssetModel).filter(AssetModel.id == asset_id).first()
if model is None:
return None
return self._to_domain(model)
def find_by_ids(self, asset_ids: list[str]) -> list[Asset]:
"""批量查询素材(单次 SQL IN 查询,避免 N+1)。"""
if not asset_ids:
return []
models = self.session.query(AssetModel).filter(AssetModel.id.in_(asset_ids)).all()
return [self._to_domain(m) for m in models]
def get(self, asset_id: str) -> Asset | None:
return self.find_by_id(asset_id)
def create(self, asset: Asset) -> Asset:
now = datetime.now(timezone.utc)
model = AssetModel(
id=asset.id,
project_id=asset.project_id,
asset_library_id=asset.library_id,
name=asset.name,
file_type=(asset.mime_type.split("/")[0] if "/" in asset.mime_type else asset.mime_type),
file_size=asset.file_size,
file_url=asset.storage_key,
storage_key=asset.storage_key,
thumbnail_url=asset.thumbnail_url,
duration=asset.duration,
width=asset.width,
height=asset.height,
fps=asset.fps,
codec=asset.codec,
status=asset.status.value,
classification_status=asset.classification_status.value,
classification_result=(json.dumps(asset.metadata) if asset.metadata else None),
quality_score=asset.quality_score,
uploaded_by_user_id=asset.uploaded_by_user_id or "system",
file_hash=asset.file_hash or None,
client_upload_id=asset.client_upload_id or None,
created_at=asset.created_at,
updated_at=now,
)
self.session.add(model)
self.session.flush()
self._sync_asset_tags(asset.id, asset.tag_ids)
self.session.commit()
return asset
def update(self, asset: Asset) -> Asset:
model = self.session.query(AssetModel).filter(AssetModel.id == asset.id).first()
if model is None:
raise ValueError(f"Asset {asset.id} not found")
model.name = asset.name
model.file_size = asset.file_size
model.file_url = asset.storage_key
model.storage_key = asset.storage_key
model.thumbnail_url = asset.thumbnail_url
model.duration = asset.duration
model.width = asset.width
model.height = asset.height
model.fps = asset.fps
model.codec = asset.codec
model.status = asset.status.value
model.classification_status = asset.classification_status.value
model.classification_result = json.dumps(asset.metadata) if asset.metadata else None
model.quality_score = asset.quality_score
model.uploaded_by_user_id = asset.uploaded_by_user_id or model.uploaded_by_user_id
model.file_hash = asset.file_hash or model.file_hash
if getattr(model, "client_upload_id", None) is None and asset.client_upload_id:
model.client_upload_id = asset.client_upload_id
model.updated_at = datetime.now(timezone.utc)
self.session.flush()
self._sync_asset_tags(asset.id, asset.tag_ids)
self.session.commit()
return asset
def delete(self, asset_id: str) -> bool:
model = self.session.query(AssetModel).filter(AssetModel.id == asset_id).first()
if model:
self.session.delete(model)
self.session.commit()
return True
return False
def batch_delete(self, asset_ids: list[str]) -> int:
"""批量删除素材(软删除,标记 status=deleted),返回实际影响数量。"""
if not asset_ids:
return 0
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
count = (
self.session.query(AssetModel)
.filter(AssetModel.id.in_(asset_ids), AssetModel.status != "deleted")
.update({AssetModel.status: "deleted", AssetModel.updated_at: now}, synchronize_session=False)
)
self.session.commit()
return count
def batch_update_metadata(self, asset_ids: list[str], metadata_patch: dict[str, object]) -> int:
"""批量更新素材 metadata(合并 patch),返回实际影响数量。"""
if not asset_ids:
return 0
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
# 逐条读取 + 合并 + 更新,保证 JSON 合并正确
models = self.session.query(AssetModel).filter(AssetModel.id.in_(asset_ids)).all()
count = 0
for model in models:
existing = {}
if model.classification_result:
try:
existing = json.loads(model.classification_result)
except Exception:
existing = {}
merged = {**existing, **metadata_patch}
model.classification_result = json.dumps(merged, ensure_ascii=False)
model.updated_at = now
count += 1
self.session.commit()
return count
def batch_add_tags(self, asset_ids: list[str], tag_ids: list[str]) -> int:
"""批量给素材添加标签(合并去重),返回实际影响数量。"""
if not asset_ids or not tag_ids:
return 0
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
clean_tag_ids = list(set(tag_ids))
count = 0
for aid in asset_ids:
# 查询现有标签
existing = {
row.tag_id
for row in self.session.query(AssetTagModel.tag_id).filter(AssetTagModel.asset_id == aid).all()
}
new_tags = [t for t in clean_tag_ids if t not in existing]
if new_tags:
for tid in new_tags:
self.session.add(AssetTagModel(asset_id=aid, tag_id=tid))
# 更新 updated_at
self.session.query(AssetModel).filter(AssetModel.id == aid).update(
{AssetModel.updated_at: now}, synchronize_session=False
)
count += 1
self.session.commit()
return count
def batch_replace_tags(self, asset_ids: list[str], tag_ids: list[str]) -> int:
"""批量替换素材标签(全量覆盖),返回实际影响数量。"""
if not asset_ids:
return 0
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
clean_tag_ids = list(set(tag_ids))
count = 0
for aid in asset_ids:
# 先删再加
self.session.query(AssetTagModel).filter(AssetTagModel.asset_id == aid).delete(synchronize_session=False)
for tid in clean_tag_ids:
self.session.add(AssetTagModel(asset_id=aid, tag_id=tid))
# 更新 updated_at
self.session.query(AssetModel).filter(AssetModel.id == aid).update(
{AssetModel.updated_at: now}, synchronize_session=False
)
count += 1
self.session.commit()
return count
def count_by_project(self, project_id: str, status: list[str] | None = None) -> int:
query = self.session.query(AssetModel).filter(AssetModel.project_id == project_id)
if status:
query = query.filter(AssetModel.status.in_(status))
return query.count()
def count_by_project_ids(self, project_ids: list[str], status: list[str] | None = None) -> int:
if not project_ids:
return 0
query = self.session.query(AssetModel).filter(AssetModel.project_id.in_(project_ids))
if status:
query = query.filter(AssetModel.status.in_(status))
return query.count()
def sum_storage_by_project_ids(self, project_ids: list[str]) -> int:
if not project_ids:
return 0
from sqlalchemy import func
result = (
self.session.query(func.coalesce(func.sum(AssetModel.file_size), 0))
.filter(AssetModel.project_id.in_(project_ids))
.scalar()
)
return int(result or 0)
def find_ready_videos_by_user(
self,
user_id: str,
*,
limit: int = 50,
) -> list[Asset]:
"""查找用户上传的所有就绪视频素材。"""
query = self.session.query(AssetModel).filter(
AssetModel.uploaded_by_user_id == user_id,
AssetModel.status == "ready",
AssetModel.file_type == "video",
)
query = query.order_by(AssetModel.created_at.desc())
if limit > 0:
query = query.limit(limit)
models = query.all()
return [self._to_domain(m) for m in models]
def search_candidates(
self,
project_id: str,
*,
file_type: str | None = None,
min_quality_score: float | None = None,
min_duration: float | None = None,
max_duration: float | None = None,
classification_category: str | None = None,
tags: list[str] | None = None,
status: str | None = None,
limit: int = 50,
) -> list[Asset]:
"""按筛选条件搜索候选素材,按质量分降序排列。"""
query = self.session.query(AssetModel).filter(
AssetModel.project_id == project_id,
)
if file_type is not None:
query = query.filter(AssetModel.file_type == file_type)
if min_quality_score is not None:
query = query.filter(AssetModel.quality_score >= min_quality_score)
if min_duration is not None:
query = query.filter(AssetModel.duration >= min_duration)
if max_duration is not None:
query = query.filter(AssetModel.duration <= max_duration)
if status is not None:
query = query.filter(AssetModel.status == status)
if classification_category is not None:
# classification_result 是 JSON Text,用 LIKE 匹配 category 字段
query = query.filter(AssetModel.classification_result.like(f'%"{classification_category}"%'))
query = query.order_by(AssetModel.quality_score.desc().nullslast())
if limit > 0:
query = query.limit(limit)
models = query.all()
candidates = [self._to_domain(m) for m in models]
# 内存中过滤 tags(tags 存在 metadata 中)
if tags:
tag_set = set(tags)
candidates = [a for a in candidates if tag_set.issubset(set(a.metadata.get("tags", [])))]
return candidates
def _to_domain(self, model: AssetModel) -> Asset:
metadata = {}
if model.classification_result:
try:
metadata = json.loads(model.classification_result)
except Exception:
metadata = {}
mime_type = model.file_type
if "/" not in mime_type:
mime_type = {
"video": "video/mp4",
"audio": "audio/mpeg",
"image": "image/jpeg",
}.get(mime_type, mime_type)
# 查询关联的 tag_ids
tag_ids = [
row.tag_id
for row in self.session.query(AssetTagModel.tag_id).filter(AssetTagModel.asset_id == model.id).all()
]
return Asset(
id=model.id,
project_id=model.project_id,
library_id=model.asset_library_id,
name=model.name,
storage_key=model.storage_key or model.file_url,
mime_type=mime_type,
file_size=int(model.file_size or 0),
thumbnail_url=model.thumbnail_url,
duration=model.duration,
width=int(model.width) if model.width is not None else None,
height=int(model.height) if model.height is not None else None,
fps=model.fps,
codec=model.codec,
status=AssetStatus(model.status),
classification_status=ClassificationStatus(model.classification_status),
quality_score=model.quality_score,
uploaded_by_user_id=model.uploaded_by_user_id,
file_hash=model.file_hash or "",
client_upload_id=getattr(model, "client_upload_id", None) or "",
metadata=metadata,
tag_ids=tag_ids,
created_at=model.created_at,
updated_at=model.updated_at,
)
def _sync_asset_tags(self, asset_id: str, tag_ids: list[str]) -> None:
"""同步素材-标签关联表(全量替换)。"""
self.session.query(AssetTagModel).filter(AssetTagModel.asset_id == asset_id).delete(synchronize_session=False)
for tag_id in tag_ids:
self.session.add(AssetTagModel(asset_id=asset_id, tag_id=tag_id))
def find_by_tag_ids(
self,
tag_ids: list[str],
skip: int = 0,
limit: int = 100,
) -> list[Asset]:
"""查找包含所有指定标签的素材。"""
if not tag_ids:
return []
from sqlalchemy import func
# 找出同时拥有所有指定 tag_id 的 asset_id
tag_set = set(tag_ids)
asset_ids = (
self.session.query(AssetTagModel.asset_id)
.filter(AssetTagModel.tag_id.in_(tag_set))
.group_by(AssetTagModel.asset_id)
.having(func.count(AssetTagModel.tag_id) == len(tag_set))
.all()
)
ids = [row[0] for row in asset_ids]
if not ids:
return []
models = self.session.query(AssetModel).filter(AssetModel.id.in_(ids)).offset(skip).limit(limit).all()
return [self._to_domain(m) for m in models]
def find_by_storage_key(self, storage_key: str) -> Asset | None:
"""按 storage_key(对应 DB 中的 file_url)查找素材。"""
model = self.session.query(AssetModel).filter(AssetModel.file_url == storage_key).first()
if model is None:
return None
return self._to_domain(model)
def find_by_library_and_file_hash(
self,
library_id: str,
file_hash: str,
) -> Asset | None:
"""按素材库 + 文件哈希查找已有素材(去重检测)。"""
if not file_hash:
return None
model = (
self.session.query(AssetModel)
.filter(
AssetModel.asset_library_id == library_id,
AssetModel.file_hash == file_hash,
)
.first()
)
if model is None:
return None
return self._to_domain(model)
def find_by_library_and_client_upload_id(
self,
library_id: str,
client_upload_id: str,
) -> Asset | None:
"""按素材库 + 客户端幂等 token 查找已有素材(complete 幂等)。"""
if not client_upload_id:
return None
model = (
self.session.query(AssetModel)
.filter(
AssetModel.asset_library_id == library_id,
AssetModel.client_upload_id == client_upload_id,
)
.first()
)
if model is None:
return None
return self._to_domain(model)
def find_recent_active_by_library_and_name(
self,
library_id: str,
name: str,
within_minutes: int = 30,
file_size: int = 0,
) -> Asset | None:
"""兜底去重:同库 + 同文件名(+同大小)且近期仍处活动状态(uploading/processing)的素材。
用于旧客户端未传 file_hash/client_upload_id 时,防止 complete 超时重试
反复创建 PROCESSING 占位记录。只命中"活动中"的近期记录,READY 历史素材不拦。
严格模式(#1714 误杀修复):file_size 必须 > 0 且与记录大小严格一致;
file_size=0(大小未知)时直接返回 None——宁可漏判(极端情况下多建一条
占位)也不可仅凭同名 + processing 误杀内容全新的视频。
"""
from datetime import datetime, timedelta, timezone
if not name:
return None
if not file_size or file_size <= 0:
return None
cutoff = datetime.now(timezone.utc) - timedelta(minutes=within_minutes)
query = self.session.query(AssetModel).filter(
AssetModel.asset_library_id == library_id,
AssetModel.name == name,
AssetModel.status.in_([AssetStatus.UPLOADING.value, AssetStatus.PROCESSING.value]),
AssetModel.created_at >= cutoff,
AssetModel.file_size == file_size,
)
model = query.order_by(AssetModel.created_at.desc()).first()
if model is None:
return None
return self._to_domain(model)