playreel/backend/pipelines/worker.py
2026-09-08 13:00:23 +09:00

78 lines
2.5 KiB
Python

"""queued 잡을 하나씩 집어 다음 게이트까지 돌리는 배경 실행기
큐가 DB에 있어 동시 실행 수를 올리거나 컨테이너를 늘려도 파이프라인 코드는 그대로
승인이 status를 queued로 되돌리면 다음 폴링에서 이어 감
"""
import asyncio
import traceback
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from pipelines.runner import QUEUED, RUNNING, run_until_gate
from tables.task import PlayreelTask, PosterAliveTask
from utils.database import session_factory
POLL_INTERVAL = 2.0
TABLES = (PlayreelTask, PosterAliveTask)
_worker: asyncio.Task | None = None
async def claim_next(session: AsyncSession):
"""queued 잡 하나를 running으로 바꿔 가져온다
워커가 하나인 동안은 경합이 없다. 늘릴 때 SELECT ... FOR UPDATE SKIP LOCKED로 바꾼다
"""
for table in TABLES:
found = await session.scalars(
select(table).where(table.status == QUEUED)
.order_by(table.created_at).limit(1))
task = found.first()
if task is not None:
task.status = RUNNING
await session.commit()
return task
return None
async def run_once() -> bool:
async with session_factory() as session:
task = await claim_next(session)
if task is None:
return False
await run_until_gate(session, task)
return True
async def requeue_running() -> None:
"""기동 시 running으로 남은 잡을 되살린다 — 어디까지 갔는지는 state가 기억한다"""
async with session_factory() as session:
for table in TABLES:
for task in await session.scalars(select(table).where(table.status == RUNNING)):
task.status = QUEUED
await session.commit()
async def loop() -> None:
await requeue_running()
while True:
try:
if not await run_once():
await asyncio.sleep(POLL_INTERVAL)
except asyncio.CancelledError:
raise
except Exception:
# 워커는 죽지 않는다. DB가 끊긴 경우라 잠시 쉬고 다시 본다
traceback.print_exc()
await asyncio.sleep(POLL_INTERVAL)
def start() -> None:
global _worker
_worker = asyncio.create_task(loop(), name="pipeline-worker")
async def stop() -> None:
if _worker is not None:
_worker.cancel()
await asyncio.gather(_worker, return_exceptions=True)