295 lines
11 KiB
Python
295 lines
11 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""썰박스 생성 잡 라이프사이클 — 생성(선차감) · 진행 · 완료 · 실패(환불) · 고아 스윕.
|
|
|
|
**모든 함수는 자체 commit 하지 않는다.** caller 가 트랜잭션을 소유하고 마지막에 한 번
|
|
커밋한다(둘 다 성공 or 둘 다 롤백). 되돌릴 수 없는 부수효과(로컬 파일 삭제)는
|
|
`finalize_task` 가 반환한 폴더를 caller 가 **커밋 성공 후에** 정리한다.
|
|
|
|
원본과 달라진 점: Task/Content 를 한 테이블(`ssul_content`)로 합쳤으므로
|
|
`finalize` 가 INSERT 가 아니라 **UPDATE** 다 — 멱등성이 자연히 확보된다.
|
|
"""
|
|
|
|
import shutil
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.credit.services.credit_service import (
|
|
deduct_credit_for_job,
|
|
refund_credit_for_job,
|
|
)
|
|
from app.ssulbox.constants import JOB_TYPE_SSUL, ORPHAN_STATUSES, SsulTaskStatus
|
|
from app.ssulbox.models import SsulContent
|
|
# castad 와 **같은 규칙으로** 지역을 뽑는다. 통합 목록에서 한 필터가 양쪽을
|
|
# 걸러야 하므로 region 값의 형식이 일치해야 한다.
|
|
from app.utils.address_parser import extract_region_from_address
|
|
from app.utils.logger import get_logger
|
|
from config import ssulbox_settings
|
|
|
|
logger = get_logger("ssulbox")
|
|
|
|
|
|
def _job_ref(content_id: int) -> str:
|
|
"""크레딧 원장 멱등 키. (job_type, job_ref, type) 유니크의 일부"""
|
|
return str(content_id)
|
|
|
|
|
|
async def create_task(
|
|
session: AsyncSession,
|
|
*,
|
|
user_uuid: str,
|
|
scenario: str,
|
|
input_text: str,
|
|
scenes: int,
|
|
seconds: int,
|
|
store_name: Optional[str] = None,
|
|
road_address: Optional[str] = None,
|
|
address: Optional[str] = None,
|
|
) -> SsulContent:
|
|
"""행 삽입 + 크레딧 선차감을 **한 트랜잭션**으로 묶는다.
|
|
|
|
사후차감이면 차감 전에 동시 요청이 들어와 크레딧 1개로 여러 개를 만들 수 있다.
|
|
잔액 부족 시 `InsufficientCreditError` 가 전파되므로 caller 가 402 로 변환한다.
|
|
|
|
`store_name`/`address` 는 검색으로 업장을 고른 경우에만 들어온다.
|
|
**생성 시점에 넣는 이유**: 통합 목록의 store_name/region 필터가 이 값에 걸리고,
|
|
생성 중인 항목도 `내 콘텐츠`에서 업장명으로 보여야 한다. 완료를 기다리면
|
|
그 사이 목록에 이름 없는 카드가 뜬다.
|
|
주소는 region 추출에만 쓰고 저장하지 않는다(castad `project` 와 달리 상세 주소를
|
|
보관할 화면이 없다). castad `/home/crawl` 과 같이 **도로명·지번을 모두** 넘겨
|
|
도로명에서 시/군 추출이 실패하면 지번으로 재시도하게 한다.
|
|
"""
|
|
row = SsulContent(
|
|
user_uuid=user_uuid,
|
|
scenario=scenario,
|
|
input=input_text,
|
|
scenes=scenes,
|
|
seconds=seconds,
|
|
status=SsulTaskStatus.QUEUED.value,
|
|
step=0,
|
|
# store_name 은 NOT NULL 이다. 아직 모르면 빈 문자열 —
|
|
# 채우는 쪽(`set_place_info`)이 falsy 검사로 "미확정"을 판단한다.
|
|
store_name=store_name or "",
|
|
region=extract_region_from_address(road_address or None, address or None)
|
|
or None,
|
|
# castad `/home/crawl` 과 동일: 도로명 우선, 없으면 지번
|
|
detail_region_info=(road_address or address or None),
|
|
)
|
|
session.add(row)
|
|
await session.flush() # autoincrement id 확보 — 크레딧 멱등 키로 쓴다
|
|
|
|
await deduct_credit_for_job(
|
|
session=session,
|
|
user_uuid=user_uuid,
|
|
amount=ssulbox_settings.SSULBOX_CREDITS_PER_VIDEO,
|
|
job_type=JOB_TYPE_SSUL,
|
|
job_ref=_job_ref(row.id),
|
|
reason="썰박스 생성",
|
|
)
|
|
logger.info(
|
|
f"[create_task] id={row.id} user={user_uuid} scenario={scenario}"
|
|
)
|
|
return row
|
|
|
|
|
|
async def mark_running(session: AsyncSession, content_id: int) -> None:
|
|
"""잡이 실제로 시작됐을 때"""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None or row.status != SsulTaskStatus.QUEUED.value:
|
|
return
|
|
row.status = SsulTaskStatus.RUNNING.value
|
|
|
|
|
|
async def update_step(session: AsyncSession, content_id: int, step: int) -> None:
|
|
"""진행 단계 갱신.
|
|
|
|
원본은 step 을 인메모리에만 뒀지만 castad 는 SSE 대신 폴링을 쓰므로
|
|
**DB 에 영속**해야 `GET /ssul/tasks/{id}` 가 진행률을 돌려줄 수 있다.
|
|
"""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None or row.status in (
|
|
SsulTaskStatus.DONE.value,
|
|
SsulTaskStatus.ERROR.value,
|
|
):
|
|
return
|
|
row.status = SsulTaskStatus.RUNNING.value
|
|
row.step = max(row.step, step) # 뒤로 가지 않는다
|
|
|
|
|
|
async def finalize_task(
|
|
session: AsyncSession,
|
|
content_id: int,
|
|
mp4_path: Path,
|
|
*,
|
|
store_name: Optional[str] = None,
|
|
region: Optional[str] = None,
|
|
) -> tuple[Optional[SsulContent], Optional[Path]]:
|
|
"""완료 처리. 같은 행을 UPDATE 하므로 **멱등**이다.
|
|
|
|
Returns:
|
|
(행, 커밋 성공 후 정리할 로컬 job 폴더 또는 None)
|
|
폴더 삭제는 되돌릴 수 없으므로 반드시 커밋이 성공한 뒤에 한다.
|
|
"""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None:
|
|
logger.warning(f"[finalize_task] 행 없음 id={content_id}")
|
|
return None, None
|
|
if row.status == SsulTaskStatus.DONE.value:
|
|
return row, None # 이미 처리됨
|
|
|
|
job_dir = mp4_path.parent
|
|
cleanup: Optional[Path] = None
|
|
video_url: Optional[str] = None
|
|
|
|
# Blob 이 설정돼 있으면 업로드하고 로컬은 커밋 후 정리한다.
|
|
# 아니면 로컬 경로를 그대로 서빙한다.
|
|
try:
|
|
from app.ssulbox.services.blob_service import upload_ssul_video
|
|
|
|
video_url = await upload_ssul_video(mp4_path, row.user_uuid, content_id)
|
|
if video_url:
|
|
cleanup = job_dir
|
|
except Exception as e:
|
|
logger.error(
|
|
f"[finalize_task] Blob 업로드 실패 id={content_id} - {type(e).__name__}: {e}",
|
|
exc_info=True,
|
|
)
|
|
|
|
if not video_url:
|
|
# 로컬 서빙 경로 (StaticFiles 마운트 기준 상대 경로)
|
|
try:
|
|
rel = mp4_path.resolve().relative_to(
|
|
ssulbox_settings.output_path.resolve()
|
|
)
|
|
video_url = f"/ssul-videos/{rel.as_posix()}"
|
|
except ValueError:
|
|
logger.error(f"[finalize_task] output 밖 경로 id={content_id} path={mp4_path}")
|
|
|
|
row.video_url = video_url
|
|
# 생성 시점에 이미 채워진 값(검색으로 사용자가 직접 고른 업장)이 우선이다.
|
|
# 여기 들어오는 값은 생성 로그에서 뒤늦게 주워온 것이므로 덮어쓰지 않는다.
|
|
if store_name and not row.store_name:
|
|
row.store_name = store_name
|
|
if region and not row.region:
|
|
row.region = region
|
|
row.status = SsulTaskStatus.DONE.value
|
|
row.step = 4
|
|
logger.info(f"[finalize_task] DONE id={content_id} url={video_url}")
|
|
return row, cleanup
|
|
|
|
|
|
async def fail_task(
|
|
session: AsyncSession, content_id: int, error: str
|
|
) -> Optional[SsulContent]:
|
|
"""실패 처리: 환불(멱등) + status=error. 이미 터미널이면 건너뛴다."""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None:
|
|
return None
|
|
if row.status in (SsulTaskStatus.DONE.value, SsulTaskStatus.ERROR.value):
|
|
return row
|
|
|
|
if row.user_uuid: # 탈퇴로 NULL 이 된 경우 환불 대상이 없다
|
|
await refund_credit_for_job(
|
|
session=session,
|
|
user_uuid=row.user_uuid,
|
|
amount=ssulbox_settings.SSULBOX_CREDITS_PER_VIDEO,
|
|
job_type=JOB_TYPE_SSUL,
|
|
job_ref=_job_ref(content_id),
|
|
reason="썰박스 생성 실패 환불",
|
|
)
|
|
|
|
row.status = SsulTaskStatus.ERROR.value
|
|
row.error = (error or "생성 실패")[:2000]
|
|
logger.info(f"[fail_task] id={content_id} error={row.error[:80]}")
|
|
return row
|
|
|
|
|
|
async def get_place_info(
|
|
session: AsyncSession, content_id: int
|
|
) -> tuple[Optional[str], Optional[str], Optional[str]]:
|
|
"""(store_name, region, detail_region_info).
|
|
|
|
크롤링으로 채울 값이 남았는지 판단하는 데 쓴다.
|
|
"""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None:
|
|
return None, None, None
|
|
return row.store_name, row.region, row.detail_region_info
|
|
|
|
|
|
async def set_place_info(
|
|
session: AsyncSession,
|
|
content_id: int,
|
|
*,
|
|
store_name: Optional[str] = None,
|
|
region: Optional[str] = None,
|
|
detail_region_info: Optional[str] = None,
|
|
) -> Optional[SsulContent]:
|
|
"""크롤링으로 얻은 업장 정보를 채운다. **이미 있는 값은 덮지 않는다.**
|
|
|
|
사용자가 검색으로 직접 고른 값이 크롤링 추정치보다 정확하므로 우선한다.
|
|
"""
|
|
row = await session.get(SsulContent, content_id)
|
|
if row is None:
|
|
return None
|
|
if store_name and not row.store_name:
|
|
row.store_name = store_name
|
|
if region and not row.region:
|
|
row.region = region
|
|
if detail_region_info and not row.detail_region_info:
|
|
row.detail_region_info = detail_region_info
|
|
logger.info(
|
|
f"[set_place_info] id={content_id} store={row.store_name!r} "
|
|
f"region={row.region!r} detail={(row.detail_region_info or '')[:30]!r}"
|
|
)
|
|
return row
|
|
|
|
|
|
async def sweep_orphans(session: AsyncSession) -> int:
|
|
"""기동 시 고아 잡 정리 — 환불 + error.
|
|
|
|
불변식: 프로세스 기동 직후 인메모리 잡은 0개이므로 DB 의 queued/running 은
|
|
**전부 이전 프로세스의 고아**다. finalize 가 단일 트랜잭션이라
|
|
"영상은 만들어졌는데 running" 같은 중간 상태는 존재하지 않는다.
|
|
|
|
⚠️ 이 불변식은 **단일 워커 전제**다. `--workers` 를 늘리면 워커 B 가 기동하며
|
|
워커 A 가 지금 돌리는 잡을 고아로 오판해 환불·error 처리한다.
|
|
"""
|
|
rows = (
|
|
(
|
|
await session.execute(
|
|
select(SsulContent).where(SsulContent.status.in_(ORPHAN_STATUSES))
|
|
)
|
|
)
|
|
.scalars()
|
|
.all()
|
|
)
|
|
for row in rows:
|
|
if row.user_uuid:
|
|
await refund_credit_for_job(
|
|
session=session,
|
|
user_uuid=row.user_uuid,
|
|
amount=ssulbox_settings.SSULBOX_CREDITS_PER_VIDEO,
|
|
job_type=JOB_TYPE_SSUL,
|
|
job_ref=_job_ref(row.id),
|
|
reason="서버 재시작 환불",
|
|
)
|
|
row.status = SsulTaskStatus.ERROR.value
|
|
row.error = "서버 재시작으로 중단되었습니다. 크레딧은 환불되었습니다."
|
|
if rows:
|
|
logger.info(f"[sweep_orphans] 고아 {len(rows)}건 환불·정리")
|
|
return len(rows)
|
|
|
|
|
|
def cleanup_job_dir(job_dir: Optional[Path]) -> None:
|
|
"""생성 산출물 폴더 삭제. **커밋이 성공한 뒤에만** 호출할 것."""
|
|
if not job_dir or not job_dir.exists():
|
|
return
|
|
try:
|
|
shutil.rmtree(job_dir)
|
|
logger.info(f"[cleanup_job_dir] 삭제 {job_dir}")
|
|
except Exception as e:
|
|
logger.warning(f"[cleanup_job_dir] 삭제 실패 {job_dir} - {e}")
|