"""alert_outbox 원장 접근. services/alert_service.py 가 부른다.""" from sqlalchemy import func, select, update from common.database.model.models import alert_outbox from common.enums import AlertStatus from common.utils.gtime import GTime async def latest_unresolved(session, dedupe_key: str): """이 dedupe_key 로 아직 안 풀린(resolved_at IS NULL) 가장 최근 알림. 없으면 None. ★ send_alert 의 중복 억제와 resolve_alert 의 "지금 알람 상태인가" 판정이 **같은 질의**를 쓴다 — 따로 구현하면 두 판단이 어긋날 수 있다.""" result = await session.execute( select(alert_outbox) .where(alert_outbox.dedupe_key == dedupe_key, alert_outbox.deleted.is_(False), alert_outbox.resolved_at.is_(None)) .order_by(alert_outbox.created_at.desc()) .limit(1) ) return result.scalars().first() async def insert(session, values: dict) -> alert_outbox: row = alert_outbox(**values) session.add(row) await session.flush() return row async def due_pending(session, limit: int = 20): """★ `next_attempt_at <= func.now()` — **DB 서버의** 지금 시각과 비교한다. 파이썬에서 계산한 GTime.UTC() 와 비교하면 앱 서버와 DB 서버의 시계가 몇 십 ms 만 어긋나도(흔하다 — 별도 컨테이너) send_alert 직후 process_outbox 를 부르는 자리에서 방금 넣은 행이 안 잡힐 수 있다(실측: 로컬에서 그렇게 재현됐다). 비교를 DB 쪽 시계 하나로 통일하면 이 경합이 없다.""" result = await session.execute( select(alert_outbox) .where(alert_outbox.status == AlertStatus.PENDING.value, alert_outbox.deleted.is_(False), alert_outbox.next_attempt_at <= func.now()) .order_by(alert_outbox.next_attempt_at) .limit(limit) ) return result.scalars().all() async def mark_sent(session, alert_id) -> None: now = GTime.UTC() await session.execute( update(alert_outbox).where(alert_outbox.alert_id == alert_id) .values(status=AlertStatus.SENT.value, sent_at=now, updated_at=now) ) async def mark_retry(session, alert_id, attempts: int, next_attempt_at) -> None: await session.execute( update(alert_outbox).where(alert_outbox.alert_id == alert_id) .values(attempts=attempts, next_attempt_at=next_attempt_at, updated_at=GTime.UTC()) ) async def mark_exhausted(session, alert_id, attempts: int) -> None: """재시도 상한 소진 — 더 시도하지 않는다(사람이 outbox 를 봐야 한다).""" await session.execute( update(alert_outbox).where(alert_outbox.alert_id == alert_id) .values(status=AlertStatus.FAILED.value, attempts=attempts, updated_at=GTime.UTC()) ) async def mark_resolved(session, alert_id) -> None: now = GTime.UTC() await session.execute( update(alert_outbox).where(alert_outbox.alert_id == alert_id) .values(resolved_at=now, updated_at=now) )