o2o-castad-backend/app/ssulbox/services/task_service.py

309 lines
12 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
from app.ssulbox.services.blob_service import generate_ssul_poster, upload_ssul_video
# 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,
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,
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:
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
try:
poster_url = await generate_ssul_poster(mp4_path, row.user_uuid, content_id)
if poster_url:
row.poster_url = poster_url
except Exception as e:
logger.warning(
f"[finalize_task] 포스터 생성 실패 id={content_id} - "
f"{type(e).__name__}: {e}",
exc_info=True,
)
# 생성 시점에 이미 채워진 값(검색으로 사용자가 직접 고른 업장)이 우선이다.
# 여기 들어오는 값은 생성 로그에서 뒤늦게 주워온 것이므로 덮어쓰지 않는다.
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], Optional[str]]:
"""(store_name, region, detail_region_info, official_site_url).
크롤링으로 채울 값이 남았는지 판단하는 데 쓴다.
"""
row = await session.get(SsulContent, content_id)
if row is None:
return None, None, None, None
return row.store_name, row.region, row.detail_region_info, row.official_site_url
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,
official_site_url: Optional[str] = None,
) -> Optional[SsulContent]:
"""크롤링으로 얻은 업장 정보를 채운다. **이미 있는 값은 덮지 않는다.**
사용자가 검색으로 직접 고른 값이 크롤링 추정치보다 정확하므로 우선한다.
`official_site_url` 만 예외로 덮어쓴다 — 호출부가 플레이스 URL 을 먼저 폴백으로
넣어 두고 크롤링에 성공하면 진짜 홈페이지로 승급시키기 때문이다.
"""
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
if official_site_url:
row.official_site_url = official_site_url[:2048]
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} "
f"site={(row.official_site_url or '')[:60]!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}")