219 lines
7.8 KiB
Python
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
|