Files
video-gen/video-gen-api/app/services/home_material/service.py
T
2026-06-30 17:10:51 +08:00

694 lines
32 KiB
Python
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.
from __future__ import annotations
import asyncio
import json
from datetime import datetime, timedelta, timezone
from typing import Any
from fastapi import HTTPException, UploadFile, status
from sqlalchemy import func, or_, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.enums.home_material import (
HOME_MATERIAL_DEFAULT_CONFIG,
HomeMaterialAssetStatus,
HomeMaterialConfigKeyEnum,
HomeMaterialMediaType,
HomeMaterialPublicResponseMode,
HomeMaterialWatermarkPosition,
HomeMaterialWatermarkSizeMode,
HomeMaterialWatermarkType,
)
from app.models.base import async_session
from app.models.home_material import HomeMaterialAsset, HomeMaterialCategory, HomeMaterialWatermark
from app.models.system_config import SystemConfig
from app.schemas.home_material import (
HomeMaterialAssetListOut,
HomeMaterialAssetOut,
HomeMaterialAssetStatusOut,
HomeMaterialAssetUpdate,
HomeMaterialCategoryCreate,
HomeMaterialCategoryListOut,
HomeMaterialCategoryOut,
HomeMaterialCategoryUpdate,
HomeMaterialConfigOut,
HomeMaterialConfigUpdate,
HomeMaterialPublicCategoryListOut,
HomeMaterialPublicFlatOut,
HomeMaterialPublicGroupedOut,
HomeMaterialRegenerateWatermarkRequest,
HomeMaterialTextWatermarkPreviewRequest,
HomeMaterialTextWatermarkPreviewResponse,
HomeMaterialUploadResultOut,
HomeMaterialWatermarkConfig,
HomeMaterialWatermarkListOut,
HomeMaterialWatermarkOut,
HomeMaterialWatermarkUpdate,
)
from app.services.home_material.query import query_service
from app.services.home_material.storage import storage_service
from app.services.home_material.watermark_processor import watermark_processor
from app.utils.id_gen import generate_id
def _json_dumps(data: Any) -> str:
return json.dumps(data, ensure_ascii=False, default=str)
def _json_loads(value: str | None) -> dict[str, Any] | None:
if not value:
return None
try:
data = json.loads(value)
return data if isinstance(data, dict) else None
except Exception:
return None
def _clean_title(value: str | None) -> str | None:
value = (value or "").strip()
return value or None
def _normalize_watermark_config_dict(value: dict[str, Any] | None, fallback_watermark_id: str | None = None) -> dict[str, Any]:
"""兼容旧水印配置。旧数据没有 watermark_type 时按 image 处理。"""
data = dict(value or {})
watermark_type = data.get("watermark_type") or HomeMaterialWatermarkType.IMAGE.value
data["watermark_type"] = watermark_type
if watermark_type == HomeMaterialWatermarkType.REPEATED_TEXT.value:
data["watermark_id"] = None
return HomeMaterialWatermarkConfig(**data).model_dump(mode="json")
if fallback_watermark_id and not data.get("watermark_id"):
data["watermark_id"] = fallback_watermark_id
return HomeMaterialWatermarkConfig(**data).model_dump(mode="json")
def _config_to_out(data: dict[str, Any] | None) -> HomeMaterialConfigOut:
raw = {**HOME_MATERIAL_DEFAULT_CONFIG, **(data or {})}
return HomeMaterialConfigOut(
enabled=bool(raw.get("enabled", False)),
title=str(raw.get("title") or HOME_MATERIAL_DEFAULT_CONFIG["title"]),
subtitle=str(raw.get("subtitle") or HOME_MATERIAL_DEFAULT_CONFIG["subtitle"]),
show_original_in_admin=bool(raw.get("show_original_in_admin", True)),
)
def _asset_snapshot(asset: HomeMaterialAsset | None) -> dict[str, Any] | None:
if not asset:
return None
return {
"id": asset.id,
"category_id": asset.category_id,
"title": asset.title,
"media_type": asset.media_type,
"status": asset.status,
"original_url": asset.original_url,
"watermarked_url": asset.watermarked_url,
"cover_url": asset.cover_url,
"watermark_id": asset.watermark_id,
"watermark_config": _json_loads(asset.watermark_config_json),
"is_active": asset.is_active,
"sort_order": asset.sort_order,
"deleted_at": asset.deleted_at,
}
def _category_snapshot(category: HomeMaterialCategory | None) -> dict[str, Any] | None:
if not category:
return None
return {
"id": category.id,
"name": category.name,
"key": category.key,
"description": category.description,
"icon": category.icon,
"is_active": category.is_active,
"sort_order": category.sort_order,
"deleted_at": category.deleted_at,
}
def _watermark_snapshot(watermark: HomeMaterialWatermark | None) -> dict[str, Any] | None:
if not watermark:
return None
return {
"id": watermark.id,
"name": watermark.name,
"file_url": watermark.file_url,
"is_default": watermark.is_default,
"is_active": watermark.is_active,
"deleted_at": watermark.deleted_at,
}
async def _flush_refresh(db: AsyncSession, obj: Any) -> None:
"""
写入后立即返回 ORM 对象前必须显式刷新。
本模块模型继承 TimestampMixinupdated_at 使用 onupdate=func.now()。
AsyncSession 下 flush 后直接访问 updated_at/created_at 可能触发隐式 IO
导致 MissingGreenlet。统一 flush + refresh,避免响应组装阶段懒加载。
"""
await db.flush()
await db.refresh(obj)
class HomeMaterialService:
"""首页素材装修主业务服务。API 层只调用本服务,内部再委托 query/storage/processor。"""
async def get_config(self, db: AsyncSession) -> HomeMaterialConfigOut:
result = await db.execute(
select(SystemConfig.value)
.where(SystemConfig.key == HomeMaterialConfigKeyEnum.SHOWCASE_CONFIG.value)
.limit(1)
)
value = result.scalar_one_or_none()
if not value:
return _config_to_out(None)
return _config_to_out(_json_loads(value))
async def save_config(self, db: AsyncSession, req: HomeMaterialConfigUpdate) -> tuple[HomeMaterialConfigOut, dict[str, Any], dict[str, Any]]:
before = await self.get_config(db)
after = HomeMaterialConfigOut(**req.model_dump())
result = await db.execute(
select(SystemConfig)
.where(SystemConfig.key == HomeMaterialConfigKeyEnum.SHOWCASE_CONFIG.value)
.limit(1)
)
config = result.scalar_one_or_none()
payload = after.model_dump()
if config:
config.value = _json_dumps(payload)
config.description = "首页素材行业装修展示配置"
db.add(config)
else:
db.add(
SystemConfig(
id=generate_id(),
key=HomeMaterialConfigKeyEnum.SHOWCASE_CONFIG.value,
value=_json_dumps(payload),
description="首页素材行业装修展示配置",
)
)
return after, before.model_dump(), after.model_dump()
async def list_categories(
self,
db: AsyncSession,
*,
page: int,
page_size: int,
keyword: str | None = None,
is_active: bool | None = None,
) -> HomeMaterialCategoryListOut:
stmt = select(HomeMaterialCategory).where(HomeMaterialCategory.deleted_at.is_(None))
count_stmt = select(func.count(HomeMaterialCategory.id)).where(HomeMaterialCategory.deleted_at.is_(None))
conditions = []
if keyword:
like = f"%{keyword}%"
conditions.append(or_(HomeMaterialCategory.name.ilike(like), HomeMaterialCategory.key.ilike(like)))
if is_active is not None:
conditions.append(HomeMaterialCategory.is_active.is_(is_active))
for cond in conditions:
stmt = stmt.where(cond)
count_stmt = count_stmt.where(cond)
total = int((await db.execute(count_stmt)).scalar_one() or 0)
result = await db.execute(
stmt.order_by(HomeMaterialCategory.sort_order.asc(), HomeMaterialCategory.created_at.desc())
.offset((page - 1) * page_size)
.limit(page_size)
)
categories = result.scalars().all()
counts = await query_service.asset_counts_map(db, [c.id for c in categories])
return HomeMaterialCategoryListOut(
items=[query_service.category_to_out(c, counts.get(c.id)) for c in categories],
total=total,
)
async def create_category(self, db: AsyncSession, req: HomeMaterialCategoryCreate, admin_id: str | None) -> tuple[HomeMaterialCategoryOut, dict[str, Any]]:
await self._ensure_category_key_available(db, req.key)
category = HomeMaterialCategory(
id=generate_id(),
name=req.name,
key=req.key,
description=req.description,
icon=req.icon,
is_active=req.is_active,
sort_order=req.sort_order,
created_by=admin_id,
updated_by=admin_id,
)
db.add(category)
await _flush_refresh(db, category)
return query_service.category_to_out(category), _category_snapshot(category) or {}
async def update_category(self, db: AsyncSession, category_id: str, req: HomeMaterialCategoryUpdate, admin_id: str | None) -> tuple[HomeMaterialCategoryOut, dict[str, Any], dict[str, Any]]:
category = await self._get_category(db, category_id)
before = _category_snapshot(category) or {}
if req.key != category.key:
await self._ensure_category_key_available(db, req.key, exclude_id=category_id)
category.name = req.name
category.key = req.key
category.description = req.description
category.icon = req.icon
category.is_active = req.is_active
category.sort_order = req.sort_order
category.updated_by = admin_id
db.add(category)
await _flush_refresh(db, category)
after = _category_snapshot(category) or {}
return query_service.category_to_out(category), before, after
async def delete_category(self, db: AsyncSession, category_id: str, admin_id: str | None) -> tuple[HomeMaterialCategory, dict[str, Any]]:
category = await self._get_category(db, category_id)
count = int((await db.execute(select(func.count(HomeMaterialAsset.id)).where(HomeMaterialAsset.category_id == category_id, HomeMaterialAsset.deleted_at.is_(None)))).scalar_one() or 0)
if count > 0:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="该行业下仍有素材,请先删除素材或禁用行业")
before = _category_snapshot(category) or {}
category.deleted_at = datetime.now(timezone.utc)
category.updated_by = admin_id
db.add(category)
return category, before
async def _ensure_category_key_available(self, db: AsyncSession, key: str, exclude_id: str | None = None) -> None:
stmt = select(HomeMaterialCategory.id).where(HomeMaterialCategory.key == key, HomeMaterialCategory.deleted_at.is_(None))
if exclude_id:
stmt = stmt.where(HomeMaterialCategory.id != exclude_id)
exists = (await db.execute(stmt.limit(1))).scalar_one_or_none()
if exists:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="行业 key 已存在")
async def _get_category(self, db: AsyncSession, category_id: str, *, active_only: bool = False) -> HomeMaterialCategory:
stmt = select(HomeMaterialCategory).where(HomeMaterialCategory.id == category_id, HomeMaterialCategory.deleted_at.is_(None))
if active_only:
stmt = stmt.where(HomeMaterialCategory.is_active.is_(True))
result = await db.execute(stmt.limit(1))
category = result.scalar_one_or_none()
if not category:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="行业不存在或已删除")
return category
async def upload_watermark(self, db: AsyncSession, *, file: UploadFile, name: str | None, is_default: bool, admin_id: str | None) -> tuple[HomeMaterialWatermarkOut, dict[str, Any]]:
stored = await storage_service.save_watermark_file(file)
info = await watermark_processor.probe(stored.storage_path)
if is_default:
await db.execute(update(HomeMaterialWatermark).where(HomeMaterialWatermark.deleted_at.is_(None)).values(is_default=False))
watermark = HomeMaterialWatermark(
id=generate_id(),
name=name or stored.file_name,
file_url=stored.file_url,
storage_path=stored.storage_path,
file_name=stored.file_name,
file_size_bytes=stored.file_size_bytes,
width=info.width,
height=info.height,
is_default=is_default,
is_active=True,
created_by=admin_id,
updated_by=admin_id,
)
db.add(watermark)
await _flush_refresh(db, watermark)
return query_service.watermark_to_out(watermark), _watermark_snapshot(watermark) or {}
async def list_watermarks(self, db: AsyncSession, *, page: int, page_size: int, is_active: bool | None = None) -> HomeMaterialWatermarkListOut:
stmt = select(HomeMaterialWatermark).where(HomeMaterialWatermark.deleted_at.is_(None))
count_stmt = select(func.count(HomeMaterialWatermark.id)).where(HomeMaterialWatermark.deleted_at.is_(None))
if is_active is not None:
stmt = stmt.where(HomeMaterialWatermark.is_active.is_(is_active))
count_stmt = count_stmt.where(HomeMaterialWatermark.is_active.is_(is_active))
total = int((await db.execute(count_stmt)).scalar_one() or 0)
result = await db.execute(
stmt.order_by(HomeMaterialWatermark.is_default.desc(), HomeMaterialWatermark.created_at.desc())
.offset((page - 1) * page_size)
.limit(page_size)
)
return HomeMaterialWatermarkListOut(items=[query_service.watermark_to_out(w) for w in result.scalars().all()], total=total)
async def update_watermark(self, db: AsyncSession, watermark_id: str, req: HomeMaterialWatermarkUpdate, admin_id: str | None) -> tuple[HomeMaterialWatermarkOut, dict[str, Any], dict[str, Any]]:
watermark = await self._get_watermark(db, watermark_id)
before = _watermark_snapshot(watermark) or {}
if req.is_default:
await db.execute(update(HomeMaterialWatermark).where(HomeMaterialWatermark.deleted_at.is_(None), HomeMaterialWatermark.id != watermark_id).values(is_default=False))
watermark.name = req.name
watermark.is_default = req.is_default
watermark.is_active = req.is_active
watermark.updated_by = admin_id
db.add(watermark)
await _flush_refresh(db, watermark)
return query_service.watermark_to_out(watermark), before, _watermark_snapshot(watermark) or {}
async def delete_watermark(self, db: AsyncSession, watermark_id: str, admin_id: str | None) -> tuple[HomeMaterialWatermark, dict[str, Any]]:
watermark = await self._get_watermark(db, watermark_id)
before = _watermark_snapshot(watermark) or {}
watermark.deleted_at = datetime.now(timezone.utc)
watermark.is_active = False
watermark.is_default = False
watermark.updated_by = admin_id
db.add(watermark)
return watermark, before
async def _get_watermark(self, db: AsyncSession, watermark_id: str, *, active_only: bool = False) -> HomeMaterialWatermark:
stmt = select(HomeMaterialWatermark).where(HomeMaterialWatermark.id == watermark_id, HomeMaterialWatermark.deleted_at.is_(None))
if active_only:
stmt = stmt.where(HomeMaterialWatermark.is_active.is_(True))
result = await db.execute(stmt.limit(1))
watermark = result.scalar_one_or_none()
if not watermark:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="水印不存在或已删除")
return watermark
async def upload_asset(
self,
db: AsyncSession,
*,
category_id: str,
file: UploadFile,
media_type: HomeMaterialMediaType,
title: str | None,
watermark_id: str | None,
watermark_file: UploadFile | None,
watermark_config: HomeMaterialWatermarkConfig,
is_active: bool,
sort_order: int,
admin_id: str | None,
) -> tuple[HomeMaterialUploadResultOut, dict[str, Any]]:
await self._get_category(db, category_id, active_only=False)
clean_title = _clean_title(title)
watermark_type = watermark_config.watermark_type
final_watermark_id: str | None = None
if watermark_type == HomeMaterialWatermarkType.IMAGE:
final_watermark_id = watermark_id
if watermark_file is not None:
watermark_out, _ = await self.upload_watermark(db, file=watermark_file, name=f"临时水印-{clean_title or '首页素材'}", is_default=False, admin_id=admin_id)
final_watermark_id = watermark_out.id
if not final_watermark_id:
default_wm = await db.execute(
select(HomeMaterialWatermark)
.where(HomeMaterialWatermark.deleted_at.is_(None), HomeMaterialWatermark.is_active.is_(True), HomeMaterialWatermark.is_default.is_(True))
.limit(1)
)
wm = default_wm.scalar_one_or_none()
final_watermark_id = wm.id if wm else None
if not final_watermark_id:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="请先选择或上传水印图片")
await self._get_watermark(db, final_watermark_id, active_only=True)
elif watermark_file is not None:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="重复文字水印模式不允许上传图片水印文件")
stored = await storage_service.save_asset_file(file, media_type)
probe = await watermark_processor.probe(stored.storage_path)
if media_type == HomeMaterialMediaType.VIDEO:
max_duration = int(getattr(settings, "HOME_MATERIAL_MAX_VIDEO_DURATION_SECONDS", 300))
if probe.duration_seconds is not None and probe.duration_seconds > max_duration:
storage_service.safe_remove(stored.storage_path)
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=f"视频时长超限,最大 {max_duration} 秒")
cfg = watermark_config.model_copy(update={"watermark_id": final_watermark_id, "opacity": round(watermark_config.opacity_level / 10, 2)})
asset = HomeMaterialAsset(
id=generate_id(),
category_id=category_id,
title=clean_title,
media_type=media_type.value,
original_url=stored.file_url,
original_storage_path=stored.storage_path,
watermark_id=final_watermark_id,
watermark_config_json=_json_dumps(cfg.model_dump(mode="json")),
status=HomeMaterialAssetStatus.PROCESSING.value,
width=probe.width,
height=probe.height,
duration_seconds=probe.duration_seconds,
file_size_bytes=stored.file_size_bytes,
is_active=is_active,
sort_order=sort_order,
created_by=admin_id,
updated_by=admin_id,
)
db.add(asset)
await _flush_refresh(db, asset)
return self._upload_result(asset, message="素材已上传,水印处理中"), _asset_snapshot(asset) or {}
async def list_assets(
self,
db: AsyncSession,
*,
page: int,
page_size: int,
category_id: str | None,
media_type: HomeMaterialMediaType | str | None,
status: HomeMaterialAssetStatus | str | None,
is_active: bool | None,
keyword: str | None,
) -> HomeMaterialAssetListOut:
config = await self.get_config(db)
items, total = await query_service.list_admin_assets(
db,
page=page,
page_size=page_size,
category_id=category_id,
media_type=media_type,
status=status,
is_active=is_active,
keyword=keyword,
include_original=config.show_original_in_admin,
)
return HomeMaterialAssetListOut(items=items, total=total)
async def get_asset_detail(self, db: AsyncSession, asset_id: str) -> HomeMaterialAssetOut:
asset = await self._get_asset(db, asset_id)
category_map = await query_service.batch_categories_map(db, [asset.category_id])
watermark_map = await query_service.batch_watermarks_map(db, [asset.watermark_id])
config = await self.get_config(db)
return query_service.asset_to_out(asset, category_map=category_map, watermark_map=watermark_map, include_original=config.show_original_in_admin)
async def update_asset(self, db: AsyncSession, asset_id: str, req: HomeMaterialAssetUpdate, admin_id: str | None) -> tuple[HomeMaterialAssetOut, dict[str, Any], dict[str, Any]]:
asset = await self._get_asset(db, asset_id)
await self._get_category(db, req.category_id)
before = _asset_snapshot(asset) or {}
asset.category_id = req.category_id
asset.title = req.title
asset.is_active = req.is_active
asset.sort_order = req.sort_order
asset.updated_by = admin_id
db.add(asset)
await _flush_refresh(db, asset)
after = _asset_snapshot(asset) or {}
return await self.get_asset_detail(db, asset_id), before, after
async def prepare_regenerate(
self,
db: AsyncSession,
asset_id: str,
req: HomeMaterialRegenerateWatermarkRequest,
admin_id: str | None,
) -> tuple[HomeMaterialUploadResultOut, dict[str, Any], dict[str, Any]]:
asset = await self._get_asset(db, asset_id)
before = _asset_snapshot(asset) or {}
req_data = req.model_dump(exclude={"wait", "wait_timeout_seconds"})
if req.watermark_type == HomeMaterialWatermarkType.IMAGE:
final_watermark_id = req.watermark_id or asset.watermark_id
await self._get_watermark(db, final_watermark_id or "", active_only=True)
cfg = HomeMaterialWatermarkConfig(**req_data).model_copy(update={"watermark_id": final_watermark_id})
else:
cfg = HomeMaterialWatermarkConfig(**req_data).model_copy(update={"watermark_id": None})
asset.watermark_id = cfg.watermark_id
asset.watermark_config_json = _json_dumps(cfg.model_dump(mode="json"))
asset.status = HomeMaterialAssetStatus.PROCESSING.value
asset.error_message = None
asset.updated_by = admin_id
db.add(asset)
await _flush_refresh(db, asset)
return self._upload_result(asset, message="水印重新生成中"), before, _asset_snapshot(asset) or {}
async def delete_asset(self, db: AsyncSession, asset_id: str, admin_id: str | None) -> tuple[HomeMaterialAsset, dict[str, Any]]:
asset = await self._get_asset(db, asset_id)
before = _asset_snapshot(asset) or {}
asset.deleted_at = datetime.now(timezone.utc)
asset.updated_by = admin_id
db.add(asset)
return asset, before
async def get_asset_status(self, db: AsyncSession, asset_id: str) -> HomeMaterialAssetStatusOut:
asset = await self._get_asset(db, asset_id)
return HomeMaterialAssetStatusOut(
id=asset.id,
status=HomeMaterialAssetStatus(asset.status),
error_message=asset.error_message,
original_url=asset.original_url,
watermarked_url=asset.watermarked_url,
cover_url=asset.cover_url,
processed_at=asset.processed_at,
)
async def _get_asset(self, db: AsyncSession, asset_id: str) -> HomeMaterialAsset:
result = await db.execute(
select(HomeMaterialAsset).where(HomeMaterialAsset.id == asset_id, HomeMaterialAsset.deleted_at.is_(None)).limit(1)
)
asset = result.scalar_one_or_none()
if not asset:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="素材不存在或已删除")
return asset
def _upload_result(self, asset: HomeMaterialAsset, *, message: str) -> HomeMaterialUploadResultOut:
return HomeMaterialUploadResultOut(
id=asset.id,
category_id=asset.category_id,
title=asset.title,
media_type=HomeMaterialMediaType(asset.media_type),
status=HomeMaterialAssetStatus(asset.status),
original_url=asset.original_url,
watermarked_url=asset.watermarked_url,
cover_url=asset.cover_url,
watermark_config=_json_loads(asset.watermark_config_json),
message=message,
)
async def wait_for_asset_result(self, asset_id: str, timeout_seconds: int) -> HomeMaterialUploadResultOut:
deadline = asyncio.get_running_loop().time() + max(1, min(timeout_seconds, 60))
while True:
async with async_session() as db:
result = await db.execute(select(HomeMaterialAsset).where(HomeMaterialAsset.id == asset_id).limit(1))
asset = result.scalar_one_or_none()
if not asset:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="素材不存在")
if asset.status in (HomeMaterialAssetStatus.SUCCESS.value, HomeMaterialAssetStatus.FAILED.value):
return self._upload_result(asset, message="水印处理完成" if asset.status == HomeMaterialAssetStatus.SUCCESS.value else "水印处理失败")
if asyncio.get_running_loop().time() >= deadline:
async with async_session() as db:
result = await db.execute(select(HomeMaterialAsset).where(HomeMaterialAsset.id == asset_id).limit(1))
asset = result.scalar_one()
return self._upload_result(asset, message="水印仍在处理中,请继续轮询状态")
await asyncio.sleep(0.5)
def start_watermark_task(self, asset_id: str) -> None:
asyncio.create_task(self._run_watermark_task(asset_id))
async def _run_watermark_task(self, asset_id: str) -> None:
try:
async with async_session() as db:
asset = await self._get_asset(db, asset_id)
cfg = _normalize_watermark_config_dict(_json_loads(asset.watermark_config_json), asset.watermark_id)
watermark_path: str | None = None
if cfg.get("watermark_type") == HomeMaterialWatermarkType.IMAGE.value:
watermark = await self._get_watermark(db, cfg.get("watermark_id") or asset.watermark_id or "", active_only=True)
watermark_path = watermark.storage_path
output_path, output_url = storage_service.build_watermarked_target(asset.id, asset.media_type)
cover_path = cover_url = None
if asset.media_type == HomeMaterialMediaType.VIDEO.value:
cover_path, cover_url = storage_service.build_cover_target(asset.id)
result = await watermark_processor.apply_watermark(
media_type=HomeMaterialMediaType(asset.media_type),
source_path=asset.original_storage_path,
watermark_path=watermark_path,
output_path=output_path,
config=cfg,
cover_path=cover_path,
)
asset.watermarked_storage_path = result.output_path
asset.watermarked_url = output_url
asset.cover_storage_path = result.cover_path
asset.cover_url = cover_url if result.cover_path else None
asset.width = result.width
asset.height = result.height
asset.duration_seconds = result.duration_seconds
asset.watermarked_file_size_bytes = result.file_size_bytes
asset.status = HomeMaterialAssetStatus.SUCCESS.value
asset.error_message = None
asset.processed_at = datetime.now(timezone.utc)
db.add(asset)
await db.commit()
except Exception as exc:
async with async_session() as db:
result = await db.execute(select(HomeMaterialAsset).where(HomeMaterialAsset.id == asset_id).limit(1))
asset = result.scalar_one_or_none()
if asset:
asset.status = HomeMaterialAssetStatus.FAILED.value
asset.error_message = str(exc)[:4000]
db.add(asset)
await db.commit()
async def preview_text_watermark_layer(self, req: HomeMaterialTextWatermarkPreviewRequest) -> HomeMaterialTextWatermarkPreviewResponse:
data_url = await asyncio.to_thread(
watermark_processor.generate_repeated_text_layer_data_url,
width=req.width,
height=req.height,
text_config=req.text_watermark,
)
return HomeMaterialTextWatermarkPreviewResponse(
width=req.width,
height=req.height,
preview_layer_data_url=data_url,
)
async def mark_stale_processing_failed(self, db: AsyncSession) -> int:
minutes = int(getattr(settings, "HOME_MATERIAL_PROCESSING_STALE_MINUTES", 30))
cutoff = datetime.now(timezone.utc) - timedelta(minutes=minutes)
result = await db.execute(
update(HomeMaterialAsset)
.where(
HomeMaterialAsset.deleted_at.is_(None),
HomeMaterialAsset.status == HomeMaterialAssetStatus.PROCESSING.value,
HomeMaterialAsset.updated_at < cutoff,
)
.values(status=HomeMaterialAssetStatus.FAILED.value, error_message="处理任务超时或服务重启,请重新生成水印")
)
return int(result.rowcount or 0)
async def get_public_categories(
self,
db: AsyncSession,
*,
with_asset_count: bool,
media_type: HomeMaterialMediaType | str | None,
only_has_assets: bool,
) -> HomeMaterialPublicCategoryListOut:
config = await self.get_config(db)
if not config.enabled:
return HomeMaterialPublicCategoryListOut(enabled=False, title=config.title, subtitle=config.subtitle, items=[])
items = await query_service.list_public_categories(db, with_asset_count=with_asset_count, media_type=media_type, only_has_assets=only_has_assets)
return HomeMaterialPublicCategoryListOut(enabled=True, title=config.title, subtitle=config.subtitle, items=items)
async def get_public_home_materials(
self,
db: AsyncSession,
*,
category_id: str | None,
category_key: str | None,
category_ids: str | None,
category_keys: str | None,
media_type: HomeMaterialMediaType | str | None,
limit_per_category: int,
include_empty_categories: bool,
response_mode: HomeMaterialPublicResponseMode,
page: int,
page_size: int,
) -> HomeMaterialPublicGroupedOut | HomeMaterialPublicFlatOut:
config = await self.get_config(db)
if not config.enabled:
if response_mode == HomeMaterialPublicResponseMode.FLAT:
return HomeMaterialPublicFlatOut(enabled=False, title=config.title, subtitle=config.subtitle, response_mode=response_mode, items=[], total=0)
return HomeMaterialPublicGroupedOut(enabled=False, title=config.title, subtitle=config.subtitle, response_mode=response_mode, categories=[])
resolved_ids = await query_service.resolve_category_ids(
db,
category_id=category_id,
category_key=category_key,
category_ids=category_ids,
category_keys=category_keys,
active_only=True,
)
if response_mode == HomeMaterialPublicResponseMode.FLAT:
items, total = await query_service.list_public_flat(db, category_ids=resolved_ids, media_type=media_type, page=page, page_size=page_size)
return HomeMaterialPublicFlatOut(enabled=True, title=config.title, subtitle=config.subtitle, response_mode=response_mode, items=items, total=total)
categories = await query_service.list_public_grouped(
db,
category_ids=resolved_ids,
media_type=media_type,
limit_per_category=limit_per_category,
include_empty_categories=include_empty_categories,
)
return HomeMaterialPublicGroupedOut(enabled=True, title=config.title, subtitle=config.subtitle, response_mode=response_mode, categories=categories)
home_material_service = HomeMaterialService()