o2o-site-AEO/solution/backend/crud/job_crud.py
Mina Choi 11d30bb3d1 [chore] solution,admin,ontology: 코드 주석을 한 줄로 — 히스토리 주석 삭제
여러 줄 주석이 설명보다 경위(예전·실측·지적)를 적고 있어 읽는 사람이 결론을 찾기 어려웠다.

- ts·tsx·js·mjs·css·py 478개: 여러 줄 주석은 첫 문장 한 줄로, 과거형·날짜 문장은 삭제
- 주석 위치는 TypeScript 파서·파이썬 tokenize/ast 로 찾는다 — 문자열 안의 # · /* 는 건드리지 않는다
- eslint·ts·noqa·type: ignore 같은 지시 주석은 그대로 둔다

파이썬 275개 정리 전후 AST 동일, TS 298개 주석 뺀 토큰 동일(빈 JSX 주석 10곳만 차이).
site·frontend·admin tsc, site vitest 105 passed

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-09-28 16:05:19 +09:00

304 lines
12 KiB
Python

"""작업 큐 CRUD — PostgreSQL 을 '제대로' 큐로 쓴다."""
import json
from sqlalchemy import text
from common.database.db_session_manager import DB_SESSION_MNG
from common.enums import DBType, DBWRType, JobStatus
# 잡 적재 시 워커를 즉시 깨우는 LISTEN/NOTIFY 채널(폴링 제거).
JOB_NOTIFY_CHANNEL = "web4ai_job"
def compute_backoff(attempts: int, base: float = 5.0, cap: float = 600.0) -> float:
"""지수 백오프(초)."""
return min(cap, base * (2 ** max(0, attempts - 1)))
class JobQueue:
DB = DBType.MAIN.value
async def _tx(self, fn):
"""쓰기 트랜잭션 — 값 반환이 필요한 큐 전이 전용 진입점."""
return await DB_SESSION_MNG.execute_lambda_write(self.DB, fn)
# 적재
async def enqueue(
self,
job_type: int,
payload: dict,
priority: int = 100,
dedupe_key: str | None = None,
max_attempts: int = 3,
) -> str | None:
"""잡 적재."""
sql = text("""
INSERT INTO jobs (job_type, priority, payload, dedupe_key, max_attempts)
VALUES (:t, :p, CAST(:payload AS jsonb), :dk, :ma)
ON CONFLICT (dedupe_key) WHERE status IN (1, 2) AND dedupe_key IS NOT NULL
DO NOTHING
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {
"t": job_type, "p": priority, "payload": json.dumps(payload, ensure_ascii=False),
"dk": dedupe_key, "ma": max_attempts,
})).first()
if row:
# 커밋 시 전달됨 → LISTEN 중인 유휴 워커를 즉시 깨운다(중복 스킵 시엔 알림 안 함).
await s.execute(text("SELECT pg_notify(:ch, '')"), {"ch": JOB_NOTIFY_CHANNEL})
return str(row[0]) if row else None
return await self._tx(run)
# 원자적 claim
async def claim(self, worker_id: str, lease_sec: int = 120) -> dict | None:
"""대기 잡 1건을 원자적으로 점유."""
sql = text("""
UPDATE jobs SET
status = 2,
worker_id = :wid,
lease_until = now() + make_interval(secs => :lease),
run_started_at = now(),
attempts = attempts + 1,
updated_at = now()
WHERE job_id = (
SELECT job_id FROM jobs
WHERE status = 1 AND run_after <= now()
ORDER BY priority ASC, created_at ASC
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING job_id, job_type, payload, attempts, max_attempts, worker_id
""")
async def run(s):
row = (await s.execute(sql, {"wid": worker_id, "lease": lease_sec})).mappings().first()
if not row:
return None
d = dict(row)
d["job_id"] = str(d["job_id"])
if isinstance(d.get("payload"), str):
d["payload"] = json.loads(d["payload"])
return d
return await self._tx(run)
# 완료/실패 (소유권 가드)
async def complete(self, job_id: str, worker_id: str, result: dict | None = None) -> bool:
sql = text("""
UPDATE jobs SET status = 3, result = CAST(:result AS jsonb),
lease_until = NULL, worker_id = NULL, updated_at = now()
WHERE job_id = CAST(:id AS uuid) AND status = 2 AND worker_id = :wid
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {
"id": job_id, "wid": worker_id,
"result": json.dumps(result, ensure_ascii=False) if result is not None else None,
})).first()
return row is not None
return await self._tx(run)
async def fail(self, job_id: str, worker_id: str, error: str, backoff_sec: float = 5.0) -> int | None:
"""실패 처리."""
sql = text("""
UPDATE jobs SET
status = CASE WHEN attempts >= max_attempts THEN 4 ELSE 1 END,
run_after = CASE WHEN attempts >= max_attempts THEN run_after
ELSE now() + make_interval(secs => :backoff) END,
last_error = :err,
lease_until = NULL,
worker_id = NULL,
updated_at = now()
WHERE job_id = CAST(:id AS uuid) AND status = 2 AND worker_id = :wid
RETURNING status
""")
async def run(s):
row = (await s.execute(sql, {
"id": job_id, "wid": worker_id, "err": error[:2000], "backoff": backoff_sec,
})).first()
return int(row[0]) if row else None
return await self._tx(run)
async def fail_permanent(self, job_id: str, worker_id: str, error: str) -> bool:
"""재시도 없이 바로 DEAD."""
sql = text("""
UPDATE jobs SET status = 4, last_error = :err,
lease_until = NULL, worker_id = NULL, updated_at = now()
WHERE job_id = CAST(:id AS uuid) AND status = 2 AND worker_id = :wid
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {"id": job_id, "wid": worker_id, "err": error[:2000]})).first()
return row is not None
return await self._tx(run)
# lease 갱신(heartbeat) / 회수(reaper)
async def renew_lease(self, job_id: str, worker_id: str, lease_sec: int = 120) -> bool:
sql = text("""
UPDATE jobs SET lease_until = now() + make_interval(secs => :lease), updated_at = now()
WHERE job_id = CAST(:id AS uuid) AND worker_id = :wid AND status = 2
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {"id": job_id, "wid": worker_id, "lease": lease_sec})).first()
return row is not None
return await self._tx(run)
async def reap(self) -> list[dict]:
"""만료된 lease(워커 사망 등)의 RUNNING 잡을 회수."""
sql = text("""
UPDATE jobs SET
status = CASE WHEN attempts >= max_attempts THEN 4 ELSE 1 END,
run_after = now(),
last_error = COALESCE(last_error, '') || ' [lease-expired reclaim]',
lease_until = NULL,
worker_id = NULL,
updated_at = now()
WHERE status = 2 AND lease_until IS NOT NULL AND lease_until < now()
RETURNING job_id, job_type, status, last_error
""")
async def run(s):
rows = (await s.execute(sql)).all()
return [
{"job_id": str(r[0]), "job_type": r[1], "status": r[2], "last_error": r[3]}
for r in rows
]
return await self._tx(run)
# 단건 조회 (상태 폴링)
async def set_progress(self, job: dict, progress: dict) -> bool:
# 회수된 옛 워커가 새 시도의 진행 상태를 덮지 못하게 한다.
sql = text("""
UPDATE jobs SET progress = CAST(:progress AS jsonb), updated_at = now()
WHERE job_id = CAST(:id AS uuid) AND status = 2
AND worker_id = :wid AND attempts = :attempt
AND lease_until > now()
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {
"id": job["job_id"], "wid": job["worker_id"], "attempt": job["attempts"],
"progress": json.dumps(progress),
})).first()
return row is not None
return await self._tx(run)
async def find_latest(self, dedupe_key: str) -> dict | None:
"""복구는 완료·실패 이력도 찾는다."""
async def run(s):
row = (await s.execute(text("""
SELECT job_id, status FROM jobs WHERE dedupe_key = :dk
ORDER BY created_at DESC, job_id DESC LIMIT 1
"""), {"dk": dedupe_key})).mappings().first()
return {**row, "job_id": str(row["job_id"])} if row else None
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
async def get(self, job_id: str) -> dict | None:
"""잡 단건 조회(읽기)."""
sql = text("""
SELECT job_id, job_type, status, priority, attempts, max_attempts,
payload, result, progress, last_error, run_after, run_started_at, created_at, updated_at
FROM jobs WHERE job_id = CAST(:id AS uuid)
""")
async def run(s):
row = (await s.execute(sql, {"id": job_id})).mappings().first()
if not row:
return None
d = dict(row)
d["job_id"] = str(d["job_id"])
for key in ("payload", "result", "progress"):
if isinstance(d.get(key), str):
d[key] = json.loads(d[key])
return d
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
async def find_active(self, dedupe_key: str) -> dict | None:
"""dedupe_key 로 활성(PENDING/RUNNING) 잡을 찾는다."""
sql = text("""
SELECT job_id, job_type, status FROM jobs
WHERE dedupe_key = :dk AND status IN (1, 2)
LIMIT 1
""")
async def run(s):
row = (await s.execute(sql, {"dk": dedupe_key})).mappings().first()
if not row:
return None
d = dict(row)
d["job_id"] = str(d["job_id"])
return d
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
# 관측(관리 API/알림용)
async def counts(self) -> dict[str, int]:
"""상태별 잡 개수."""
async def run(s):
rows = (await s.execute(text("SELECT status, count(*) FROM jobs GROUP BY status"))).all()
by_val = {int(st): int(c) for st, c in rows}
return {js.name: by_val.get(js.value, 0) for js in JobStatus}
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
async def ops(self) -> dict:
"""운영 스냅샷(모니터링·알림용): 상태별 카운트 + 큐 지연(가장 오래된 PENDING 나이) + 최근 1시간 DEAD + stuck(좀비 신호)."""
sql = text("""
SELECT
count(*) FILTER (WHERE status = 1) AS pending,
count(*) FILTER (WHERE status = 2) AS running,
count(*) FILTER (WHERE status = 3) AS done,
count(*) FILTER (WHERE status = 4) AS dead,
count(*) FILTER (WHERE status = 4 AND updated_at > now() - interval '1 hour') AS dead_1h,
count(*) FILTER (WHERE status = 2 AND (
(lease_until IS NOT NULL AND lease_until < now())
OR run_started_at < now() - interval '10 minutes'
)) AS stuck_running,
COALESCE(EXTRACT(EPOCH FROM (now() - min(created_at) FILTER (WHERE status = 1)))::int, 0) AS oldest_pending_sec,
count(*) FILTER (WHERE last_error LIKE 'JobDeadlineExceeded%'
AND updated_at > now() - interval '1 hour') AS deadline_1h
FROM jobs
""")
async def run(s):
return {k: int(v) for k, v in dict((await s.execute(sql)).mappings().first()).items()}
return await DB_SESSION_MNG.execute_lambda(self.DB, DBWRType.DB_READ.value, run)
async def requeue(self, job_id: str) -> str | None:
"""DEAD 잡 재큐(관리자 액션): attempts 리셋 + PENDING 전이 + 워커 깨움."""
sql = text("""
UPDATE jobs SET status = 1, attempts = 0, run_after = now(),
lease_until = NULL, worker_id = NULL, run_started_at = NULL,
last_error = NULL, updated_at = now()
WHERE job_id = CAST(:jid AS uuid) AND status = 4
RETURNING job_id
""")
async def run(s):
row = (await s.execute(sql, {"jid": job_id})).first()
if row:
await s.execute(text("SELECT pg_notify(:ch, '')"), {"ch": JOB_NOTIFY_CHANNEL})
return str(row[0]) if row else None
return await self._tx(run)