# -*- 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}")