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

219 lines
7.8 KiB
Python

# -*- coding: utf-8 -*-
"""F1(무빙 포스터) 잡 라이프사이클 + 크레딧.
모든 함수는 **자체 commit 하지 않는다.** 트랜잭션 소유는 라우터/백그라운드 태스크다.
환불 정책 (F2 와 다르다):
F1 실패는 대부분 재시도 가능하다(제목 게이트, 모델 큐 등 — 프론트에 재시도
버튼이 있다). 실패 관측 즉시 환불하면 크레딧 원장의 (job_type, job_ref, type)
유니크 제약 때문에 재시도 시 재차감이 불가능해져 "환불받고 공짜 재시도"가 된다.
그래서 실패 시에는 환불하지 않고, **잡을 버릴 때(삭제)와 고아 스윕에서만** 환불한다.
"""
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_F1,
ORPHAN_STATUSES,
STATUS_ARCHIVING,
STATUS_DONE,
STATUS_FAILED,
)
from app.p2v.exceptions import P2vJobNotFoundError
from app.p2v.models import P2vF1Job
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: str
) -> P2vF1Job:
"""잡 행 생성 + 크레딧 선차감 (한 트랜잭션 — 호출부가 commit).
행을 먼저 flush 해 id 를 확보한 뒤 그 id 를 멱등 키(job_ref)로 차감한다.
크레딧이 부족하면 InsufficientCreditError 가 올라가고 행도 롤백된다 —
P2V 는 아직 호출되지 않았으므로 실비용(OpenAI 등)이 나가지 않는다.
Raises:
InsufficientCreditError: 잔액 부족 (전역 핸들러가 402 로 변환)
"""
row = P2vF1Job(
user_uuid=user_uuid,
name=name or None,
status="queued",
credit_amount=p2v_settings.P2V_CREDITS_PER_F1,
)
session.add(row)
await session.flush() # id 확보
await deduct_credit_for_job(
session=session,
user_uuid=user_uuid,
amount=p2v_settings.P2V_CREDITS_PER_F1,
job_type=JOB_TYPE_P2V_F1,
job_ref=str(row.id),
reason="무빙 포스터 생성",
)
logger.info(f"[f1] 잡 생성 id={row.id} user={user_uuid}")
return row
async def get_owned(
session: AsyncSession, job_id: int, user_uuid: str
) -> P2vF1Job:
"""소유권 검사 포함 조회. 남의 잡이면 존재를 알리지 않고 동일하게 404."""
row = await session.get(P2vF1Job, job_id)
if row is None or row.user_uuid != user_uuid:
raise P2vJobNotFoundError()
return row
async def refund_job(
session: AsyncSession, row: P2vF1Job, reason: str
) -> None:
"""선차감분 환불 (멱등 — 이미 환불했거나 차감 기록이 없으면 조용히 넘어간다)."""
await refund_credit_for_job(
session=session,
user_uuid=row.user_uuid,
amount=row.credit_amount,
job_type=JOB_TYPE_P2V_F1,
job_ref=str(row.id),
reason=reason,
)
async def fail_job(
session: AsyncSession, row: P2vF1Job, detail: str, *, refund: bool
) -> None:
"""실패 처리.
Args:
refund: True 면 환불까지. P2V 호출 자체가 실패했거나(실비용 없음)
잡이 P2V 에서 사라져 재시도가 불가능할 때만 True 로 준다.
일반 스테이지 실패는 재시도 가능하므로 False (모듈 docstring 참조).
"""
row.status = STATUS_FAILED
row.error = detail[:1000]
if refund:
await refund_job(session, row, reason=f"무빙 포스터 실패: {detail[:100]}")
logger.info(f"[f1] 실패 처리 id={row.id} refund={refund} detail={detail[:200]}")
def sync_from_p2v(row: P2vF1Job, p2v_job: dict) -> None:
"""P2V 잡 JSON → castad 행 미러링 (순수 상태 복사 — DB 호출 없음).
'done' 은 여기서 반영하지 않는다 — Blob 아카이빙이 끝나야 done 이다
(archive_service 가 archiving → done 전이를 소유한다).
실패도 상태만 미러하고 환불하지 않는다 (재시도 가능 — 모듈 docstring 참조).
"""
p2v_status = p2v_job.get("status", "")
if p2v_status in ("queued", "running", "awaiting_review"):
row.status = p2v_status
elif p2v_status == "failed":
row.status = STATUS_FAILED
err = p2v_job.get("error") or {}
row.error = f"{err.get('stage', '?')}: {err.get('detail', '')}"[:1000]
if p2v_job.get("name"):
row.name = str(p2v_job["name"])[:100]
if p2v_job.get("narration") is not None:
row.narration = p2v_job["narration"]
md = p2v_job.get("metadata")
if md:
row.event_name = (md.get("event_name") or "")[:200] or None
row.date_text = (md.get("date_text") or "")[:100] or None
row.place = (md.get("place") or "")[:200] or None
row.keywords = md.get("keywords")
if p2v_job.get("duration") is not None:
row.duration = p2v_job["duration"]
async def try_claim_archiving(session: AsyncSession, job_id: int) -> bool:
"""Blob 아카이빙 착수권 선점 (낙관적 잠금).
폴링이 겹쳐도 UPDATE 는 한 요청만 성공한다. rowcount 0 이면 다른 폴링이
이미 집어갔거나 이미 아카이브된 것이므로 조용히 물러난다.
"""
result = await session.execute(
update(P2vF1Job)
.where(
P2vF1Job.id == job_id,
P2vF1Job.status != STATUS_ARCHIVING,
P2vF1Job.archived_at.is_(None),
)
.values(status=STATUS_ARCHIVING)
.execution_options(synchronize_session=False)
)
return result.rowcount > 0
async def finish_archive(
session: AsyncSession,
row: P2vF1Job,
*,
video_url: Optional[str],
thumbnail_url: Optional[str],
) -> None:
"""아카이빙 완료 — Blob URL 기록 + done 전이."""
row.p2v_video_url = video_url
row.poster_url = thumbnail_url
row.archived_at = datetime.now()
row.status = STATUS_DONE
logger.info(f"[f1] 아카이브 완료 id={row.id} video={bool(video_url)}")
async def sweep_orphans(session: AsyncSession) -> int:
"""방치된 잡 정리 (앱 기동 시 lifespan 훅에서 호출).
1) 진행/검수 상태로 P2V_ORPHAN_TIMEOUT_HOURS 를 넘긴 잡 → 환불 + 실패
(검수 화면에서 이탈한 사용자의 선차감분이 증발하는 것을 막는다 — 설계 R-1)
2) archiving 으로 P2V_ARCHIVE_STALE_MINUTES 넘게 방치된 잡 → 'done'(archived_at
NULL) 으로 되돌린다. 다음 폴링이 아카이빙을 다시 선점한다 (설계 R-3 크래시 복구)
시각 비교는 DB 시계(func.now())로 한다 — 앱과 DB 의 타임존이 다를 수 있다.
"""
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(P2vF1Job).where(
P2vF1Job.status.in_(ORPHAN_STATUSES),
P2vF1Job.created_at < orphan_cutoff,
)
)
for row in result.scalars().all():
await fail_job(
session, row, "시간 초과 — 크레딧을 자동 환불했습니다", refund=True
)
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(P2vF1Job)
.where(
P2vF1Job.status == STATUS_ARCHIVING,
P2vF1Job.updated_at < stale_cutoff,
)
.values(status=STATUS_DONE)
.execution_options(synchronize_session=False)
)
if result.rowcount:
logger.info(f"[f1] 정체된 archiving {result.rowcount}건 재시도 가능으로 복구")
return swept