o2o-castad-backend/app/p2v/services/f2_service.py
2026-08-26 13:52:18 +09:00

191 lines
6.2 KiB
Python

# -*- coding: utf-8 -*-
"""F2(포스터 스타일링) 잡 라이프사이클 + 크레딧.
모든 함수는 **자체 commit 하지 않는다.** 트랜잭션 소유는 라우터/백그라운드 태스크다.
F1 과 달리 F2 에는 재시도 경로가 없다(P2V API 에 retry 엔드포인트 자체가 없음).
따라서 실패를 관측하는 즉시 환불한다.
"""
from datetime import datetime
from typing import Optional
from sqlalchemy import func, select, text, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.credit.services.credit_service import (
deduct_credit_for_job,
refund_credit_for_job,
)
from app.p2v.constants import (
JOB_TYPE_P2V_F2,
STATUS_ARCHIVING,
STATUS_DONE,
STATUS_FAILED,
)
from app.p2v.exceptions import P2vJobNotFoundError
from app.p2v.models import P2vF2Job
from app.utils.logger import get_logger
from config import p2v_settings
logger = get_logger("p2v")
async def create_job(
session: AsyncSession,
*,
user_uuid: str,
name: Optional[str],
template_id: str,
template_name: Optional[str],
license: Optional[str],
format: str,
) -> P2vF2Job:
"""잡 행 생성 + 크레딧 선차감 (한 트랜잭션 — 호출부가 commit).
템플릿 이름·라이선스는 이 시점의 스냅샷이다 — 템플릿이 나중에 삭제되거나
등급이 바뀌어도 "이 결과물이 어떤 등급 레퍼런스로 만들어졌는지"가 남는다.
Raises:
InsufficientCreditError: 잔액 부족 (전역 핸들러가 402 로 변환)
"""
row = P2vF2Job(
user_uuid=user_uuid,
name=(name or "")[:100] or None,
template_id=template_id,
template_name=(template_name or "")[:100] or None,
license=(license or "")[:20] or None,
format=format,
status="queued",
credit_amount=p2v_settings.P2V_CREDITS_PER_F2,
)
session.add(row)
await session.flush() # id 확보
await deduct_credit_for_job(
session=session,
user_uuid=user_uuid,
amount=p2v_settings.P2V_CREDITS_PER_F2,
job_type=JOB_TYPE_P2V_F2,
job_ref=str(row.id),
reason="포스터 스타일링",
)
logger.info(f"[f2] 잡 생성 id={row.id} user={user_uuid} template={template_id}")
return row
async def get_owned(
session: AsyncSession, job_id: int, user_uuid: str
) -> P2vF2Job:
"""소유권 검사 포함 조회. 남의 잡이면 존재를 알리지 않고 동일하게 404."""
row = await session.get(P2vF2Job, job_id)
if row is None or row.user_uuid != user_uuid:
raise P2vJobNotFoundError()
return row
async def refund_job(session: AsyncSession, row: P2vF2Job, reason: str) -> None:
"""선차감분 환불 (멱등 — 이미 환불했거나 차감 기록이 없으면 조용히 넘어간다)."""
await refund_credit_for_job(
session=session,
user_uuid=row.user_uuid,
amount=row.credit_amount,
job_type=JOB_TYPE_P2V_F2,
job_ref=str(row.id),
reason=reason,
)
async def fail_job(session: AsyncSession, row: P2vF2Job, detail: str) -> None:
"""실패 처리 + 즉시 환불 (F2 는 재시도 경로가 없다)."""
row.status = STATUS_FAILED
row.error = detail[:1000]
await refund_credit_for_job(
session=session,
user_uuid=row.user_uuid,
amount=row.credit_amount,
job_type=JOB_TYPE_P2V_F2,
job_ref=str(row.id),
reason=f"포스터 스타일링 실패: {detail[:100]}",
)
logger.info(f"[f2] 실패 처리(환불) id={row.id} detail={detail[:200]}")
async def sync_from_p2v(
session: AsyncSession, row: P2vF2Job, p2v_job: dict
) -> None:
"""P2V 잡 JSON → castad 행 미러링.
'done' 은 반영하지 않는다 — Blob 아카이빙 완료가 done 이다.
'failed' 는 즉시 환불한다 (F1 과 다른 지점).
"""
p2v_status = p2v_job.get("status", "")
if p2v_status in ("queued", "running"):
row.status = p2v_status
elif p2v_status == "failed" and row.status != STATUS_FAILED:
err = p2v_job.get("error") or {}
await fail_job(
session, row, f"{err.get('stage', '?')}: {err.get('detail', '')}"
)
async def try_claim_archiving(session: AsyncSession, job_id: int) -> bool:
"""Blob 아카이빙 착수권 선점 (낙관적 잠금 — f1_service 와 동일 방식)."""
result = await session.execute(
update(P2vF2Job)
.where(
P2vF2Job.id == job_id,
P2vF2Job.status != STATUS_ARCHIVING,
P2vF2Job.archived_at.is_(None),
)
.values(status=STATUS_ARCHIVING)
.execution_options(synchronize_session=False)
)
return result.rowcount > 0
async def finish_archive(
session: AsyncSession, row: P2vF2Job, *, image_url: Optional[str]
) -> None:
"""아카이빙 완료 — Blob URL 기록 + done 전이."""
row.p2v_poster_url = image_url
row.archived_at = datetime.now()
row.status = STATUS_DONE
logger.info(f"[f2] 아카이브 완료 id={row.id} image={bool(image_url)}")
async def sweep_orphans(session: AsyncSession) -> int:
"""방치된 잡 정리 — f1_service.sweep_orphans 와 같은 규칙 (검수 단계만 없다)."""
swept = 0
orphan_cutoff = func.date_sub(
func.now(), text(f"INTERVAL {int(p2v_settings.P2V_ORPHAN_TIMEOUT_HOURS)} HOUR")
)
result = await session.execute(
select(P2vF2Job).where(
P2vF2Job.status.in_(("queued", "running")),
P2vF2Job.created_at < orphan_cutoff,
)
)
for row in result.scalars().all():
await fail_job(session, row, "시간 초과 — 크레딧을 자동 환불했습니다")
swept += 1
stale_cutoff = func.date_sub(
func.now(),
text(f"INTERVAL {int(p2v_settings.P2V_ARCHIVE_STALE_MINUTES)} MINUTE"),
)
result = await session.execute(
update(P2vF2Job)
.where(
P2vF2Job.status == STATUS_ARCHIVING,
P2vF2Job.updated_at < stale_cutoff,
)
.values(status=STATUS_DONE)
.execution_options(synchronize_session=False)
)
if result.rowcount:
logger.info(f"[f2] 정체된 archiving {result.rowcount}건 재시도 가능으로 복구")
return swept