o2o-triple-pick/backend/app/worker.py

518 lines
22 KiB
Python

"""백그라운드 워커 — APScheduler (별도 컨테이너).
조사된 '필요한 시간들'을 모두 자동 처리하는 단일 프로세스:
1) 상태 전이 (status_tick_seconds 주기, 기본 60초)
scheduled → open(킥오프 D-2) → locked(킥오프 정각) — now 기준 자동 갱신.
2) AI 예측 생성 (매일 KST ai_generate_hour:ai_generate_minute, 기본 00:05)
남은(미종료) 경기에 대해 GPT/Claude/Gemini 실 API 호출 → ai_predictions 갱신.
모델별 독립 처리(키 없거나 실패해도 나머지 모델은 진행).
3) 결과 메일 발송 (status_tick 와 함께 점검)
경기 종료(finished_at) 후 result_email_delay_minutes(기본 180=3시간) 경과 시,
구독자(notify=True+email)에게 개인화 결과 메일 발송 → notified 표시.
실행: python -m app.worker
"""
from __future__ import annotations
import asyncio
import logging
from datetime import timedelta
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from sqlalchemy import select
from sqlalchemy.orm import selectinload
from .config import settings
from .database import SessionLocal, init_db
from .domain import compute_phase, ensure_aware, now_utc
from .models import AIPrediction, Match, UserPrediction
from .scoring import load_scoring_data
from .services import baseball_data, baseball_details, football_data
from .services.ai import MatchContext, PROVIDERS, ProviderUnavailable
from .services.baseball_fetch import fetch_baseball_results, fetch_baseball_schedule
from .services.baseball_sync import match_seq, sync_baseball_schedule
from .services.email import EmailUnavailable, build_result_email, send_email
from .services.grading import apply_result
from .services.schedule_fetch import fetch_results, fetch_schedule
from .services.schedule_sync import sync_schedule
logging.basicConfig(level=logging.INFO)
log = logging.getLogger("triplepick.worker")
# ── 0) 경기 일정 동기화 (외부 크롤링) ──────────────────────
async def sync_schedule_job() -> None:
"""월드컵(축구) 일정 동기화 — 기존 로직 그대로."""
if "wc" not in settings.league_list:
return
records = await fetch_schedule()
if not records:
return
async with SessionLocal() as db:
result = await sync_schedule(db, records)
# 일정 변경 후 즉시 상태/투표시간 재평가
await tick_status()
# 새 경기가 삽입됐으면 그 경기 AI 예측을 바로 생성(다음 주기까지 비어있지 않게)
if result.get("inserted"):
log.info("schedule: 신규 %d경기 → AI 예측 즉시 생성", result["inserted"])
await generate_ai_predictions()
async def sync_baseball_job() -> None:
"""야구(kbo/mlb) 일정 동기화 + 프리뷰·순위 캐시 갱신."""
inserted_any = False
for league in settings.league_list:
if league not in ("kbo", "mlb"):
continue
records = await fetch_baseball_schedule(league)
if not records:
continue
async with SessionLocal() as db:
result = await sync_baseball_schedule(db, league, records)
rows = (
await db.execute(
select(Match).where(
Match.league == league, Match.result_outcome.is_(None)
)
)
).scalars().all()
try:
await baseball_details.refresh_baseball_details(db, league, rows)
except Exception as e: # noqa: BLE001
log.warning("baseball details(%s) 갱신 실패: %s", league, e)
inserted_any = inserted_any or bool(result.get("inserted"))
await tick_status()
if inserted_any:
await generate_ai_predictions()
# ── 1) 상태 전이 ────────────────────────────────────────────
async def tick_status() -> None:
async with SessionLocal() as db:
now = now_utc()
matches = (await db.execute(select(Match))).scalars().all()
changed = 0
for m in matches:
# 화면용 phase 와 동일 기준으로 status 저장(단일 기준):
# scheduled → open → locked(투표종료) → live(경기중) → finished
target = compute_phase(m, now)
if m.status != target:
m.status = target
changed += 1
if changed:
await db.commit()
log.info("status_tick: %d matches updated", changed)
await maybe_send_result_emails()
# ── 2) AI 예측 생성 (실연동) ────────────────────────────────
async def generate_ai_predictions(only_missing: bool = True) -> None:
"""미종료 경기의 AI 예측 생성.
only_missing=True(기본): 이미 있는 (경기×모델) 예측은 보존하고 **비어 있는 것만**
생성한다(직전에 실패/누락된 모델만 채워짐 — API 비용↓, 예측 안정).
only_missing=False: 전부 재생성(덮어쓰기) — 수동 재생성 스크립트용.
"""
if not settings.ai_enabled:
log.info("AI 예측 생성 스킵 — AI_ENABLED=false (로컬 테스트 모드)")
return
async with SessionLocal() as db:
# 킥오프가 lookahead(기본 22h) 이내로 다가온 미시작 경기만 생성.
# 야구 선발투수 예고(전날 저녁)가 나온 뒤 시점이라 데이터 품질이 좋다.
# 킥오프가 지난 경기는 제외(경기 중 생성 = 사후 예측 방지). 취소 경기 제외.
now = now_utc()
horizon = now + timedelta(hours=settings.ai_generate_lookahead_hours)
rows = (
await db.execute(
select(Match)
.where(Match.result_outcome.is_(None))
.options(selectinload(Match.predictions))
)
).scalars().all()
matches = [
m for m in rows
if m.status != "cancelled"
and now < ensure_aware(m.kickoff_at) <= horizon
]
for m in matches:
# 실데이터 블록 — 리그별 소스(축구=API-Football 캐시, 야구=자체DB+프리뷰).
if m.league in ("kbo", "mlb"):
data_block = await baseball_data.build_baseball_data_block(db, m)
else:
data_block = await football_data.build_data_block(db, m)
ctx = MatchContext(
team_a=m.team_a_name,
team_b=m.team_b_name,
venue=m.venue,
kickoff=m.kickoff_at.isoformat(),
data_block=data_block,
league=m.league,
)
for model, fn in PROVIDERS.items():
pred = next((p for p in m.predictions if p.model == model), None)
if only_missing and pred is not None:
continue # 이미 작성됨 — 보존, API 호출 안 함
try:
data = await fn(ctx)
except ProviderUnavailable as e:
log.warning("AI %s skipped (%s)", model, e)
continue
except Exception as e: # noqa: BLE001
log.error("AI %s failed for %s: %s", model, m.match_id, e)
continue
if pred is None: # 신규 — 위에서 못 찾았으면 생성
pred = AIPrediction(match_id=m.match_id, model=model)
db.add(pred)
m.predictions.append(pred)
pred.outcome = data["outcome"]
pred.score_a = data["scoreA"]
pred.score_b = data["scoreB"]
pred.confidence_pct = data["confidencePct"]
pred.reason_ko = data["reasonKo"]
pred.reason_en = data["reasonEn"]
pred.generated_at = now_utc()
pred.source = "llm"
log.info("AI %s%s %d-%d", model, m.match_id, data["scoreA"], data["scoreB"])
# 경기 단위 커밋 — 수백 콜 도중 재시작/오류가 나도 진행분을 잃지
# 않는다(전체 커밋이면 중단 시 API 비용을 다시 지출하게 됨).
await db.commit()
# ── 2.1) 축구 데이터 수집 (예측 생성 전에 캐시를 채워둔다) ──
async def refresh_football_data() -> None:
"""전 팀의 팀/H2H 축구 데이터를 캐시에 한 번씩 선수집.
예측이 데이터를 쓰려면 미리 캐시에 있어야 하므로, 임박 경기만 기다리지 않고
모든 미종료 경기의 팀을 대상으로 한다. 팀당 1회만 받고(fetch-once), 무료 한도
(100/일)를 넘지 않게 하루 호출 예산만큼만 받아 며칠에 걸쳐 누적한다. 전 팀이
캐시되면 이후 잡은 호출 0(자동 무동작). 키 미설정이면 no-op.
"""
if not football_data.enabled():
return
async with SessionLocal() as db:
rows = (
await db.execute(
select(Match).where(
Match.result_outcome.is_(None), Match.league == "wc"
)
)
).scalars().all()
# 가까운 경기부터 우선 수집(하루 예산 소진 시 먼 경기 팀은 다음 잡에서).
matches = sorted(rows, key=lambda m: ensure_aware(m.kickoff_at))
await football_data.refresh(db, matches)
# ── 2.5) 결과 자동 정산 (관리자 입력 불필요) ────────────────
async def settle_matches() -> None:
"""킥오프가 지난 미정산 경기를, 소스가 FINISHED 로 보고하는 즉시 채점·종료.
또한 최근 N일(result_recheck_days) 내 종료된 경기는 외부 소스와 대조해, 소스가
스코어를 정정하면(잠정값→확정값 등) 자동 갱신·재채점한다.
킥오프 이후의 미정산 경기를 폴링하되, 외부 소스(football-data)가 FINISHED 로
확정 스코어를 줄 때만 apply_result → 채점·유저 포인트 집계·finished_at·메일.
경기 중(IN_PLAY)에는 FINISHED 가 아니라 결과에 안 잡혀 자동 스킵된다(조기 종료 없음).
아직 FINISHED 가 아니면 다음 틱(5분)에 재시도 → 종료 후 최대 1틱 내 반영.
재확인은 '최근 N일 종료 경기'로만 한정 → 오래된 경기는 대상에서 빠져 부하 bounded.
소스 호출은 신규/재확인이 모두 같은 단일 fetch_results() 응답을 공유한다.
"""
now = now_utc()
recheck_cutoff = (
now - timedelta(days=settings.result_recheck_days)
if settings.result_recheck_days > 0
else None
)
async with SessionLocal() as db:
pending = (
await db.execute(
select(Match).where(
Match.league == "wc", Match.result_outcome.is_(None)
)
)
).scalars().all()
# 킥오프가 지난 경기만 (소스가 FINISHED 줄 수 있는 시점)
due = [m.match_id for m in pending if ensure_aware(m.kickoff_at) <= now]
# 최근 종료 경기(소스 정정 반영용) — TZ 안전하게 파이썬에서 기간 필터.
recheck_ids: list[str] = []
if recheck_cutoff is not None:
finished = (
await db.execute(
select(Match).where(
Match.league == "wc",
Match.result_outcome.is_not(None),
Match.finished_at.is_not(None),
)
)
).scalars().all()
recheck_ids = [
m.match_id
for m in finished
if m.finished_at and ensure_aware(m.finished_at) >= recheck_cutoff
]
if due or recheck_ids:
results = await fetch_results()
# 정방향/역방향 키 모두 등록 — 소스의 홈/원정 순서가 우리와 달라도 매칭.
# 라운드(round_label)를 키에 포함해 같은 두 팀의 조별리그·토너먼트 경기를 구분.
# round_label 없는 소스(openfootball) 대비 라운드 무시 키도 함께 등록(폴백).
by_pair: dict[tuple, tuple[int, int]] = {}
for r in results:
rl = r.get("roundLabel", "")
a2, b2, sa, sb = r["teamA"], r["teamB"], r["scoreA"], r["scoreB"]
for key in ((rl, a2, b2), (a2, b2)):
by_pair[key] = (sa, sb)
for key in ((rl, b2, a2), (b2, a2)): # 뒤집어 저장
by_pair[key] = (sb, sa)
def _lookup(m: Match) -> tuple[int, int] | None:
# 라운드까지 일치하는 결과 우선, 없으면 팀쌍만으로 폴백.
return by_pair.get((m.round_label, m.team_a_code, m.team_b_code)) or by_pair.get(
(m.team_a_code, m.team_b_code)
)
settled = corrected = 0
async with SessionLocal() as db:
# 1) 신규 정산 — 미정산 경기에 FINISHED 스코어 반영
for mid in due:
m = await db.get(Match, mid)
if m is None or m.result_outcome is not None:
continue
sc = _lookup(m)
if not sc:
continue
await apply_result(db, m, sc[0], sc[1]) # 우리 팀A/팀B 기준으로 채점+집계+commit
settled += 1
log.info("settle: %s 자동 정산 %d-%d", mid, sc[0], sc[1])
# 2) 최근 종료 경기 재확인 — 소스 스코어가 DB와 다르면 정정+재채점
for mid in recheck_ids:
m = await db.get(Match, mid)
if m is None or m.result_outcome is None:
continue
sc = _lookup(m)
if not sc:
continue # 소스에 아직 없으면 기존값 유지
if (m.result_score_a, m.result_score_b) == (sc[0], sc[1]):
continue # 동일 — 변경 없음(대부분 여기서 종료, 재채점 안 함)
old_a, old_b = m.result_score_a, m.result_score_b
await apply_result(db, m, sc[0], sc[1]) # 정정+재채점+집계+commit
corrected += 1
log.warning(
"settle: %s 결과 정정 %s-%s%d-%d (소스 변경 반영)",
mid, old_a, old_b, sc[0], sc[1],
)
if not settled and not corrected:
log.info(
"settle: 신규 %d·재확인 %d경기 — 변경 없음", len(due), len(recheck_ids)
)
await maybe_send_result_emails()
# ── 2.6) 야구 결과 자동 정산 — (리그, KST 날짜, 팀쌍) 키 매칭 ──
async def settle_baseball() -> None:
_KST = timedelta(hours=9)
now = now_utc()
for league in settings.league_list:
if league not in ("kbo", "mlb"):
continue
async with SessionLocal() as db:
pending = (
await db.execute(
select(Match).where(
Match.league == league,
Match.result_outcome.is_(None),
Match.status != "cancelled",
)
)
).scalars().all()
due = [m.match_id for m in pending if ensure_aware(m.kickoff_at) <= now]
if not due:
continue
results = await fetch_baseball_results(league)
by_key: dict[tuple, tuple[int, int]] = {}
for r in results:
d, a2, b2, s = r["dateKst"], r["teamA"], r["teamB"], r.get("seq", 1)
by_key[(d, a2, b2, s)] = (r["scoreA"], r["scoreB"])
by_key[(d, b2, a2, s)] = (r["scoreB"], r["scoreA"])
settled = 0
async with SessionLocal() as db:
for mid in due:
m = await db.get(Match, mid)
if m is None or m.result_outcome is not None:
continue
d = (ensure_aware(m.kickoff_at) + _KST).strftime("%Y%m%d")
sc = by_key.get((d, m.team_a_code, m.team_b_code, match_seq(mid)))
if not sc:
continue
await apply_result(db, m, sc[0], sc[1])
settled += 1
log.info("settle(%s): %s 자동 정산 %d-%d", league, mid, sc[0], sc[1])
if settled:
await maybe_send_result_emails()
# ── 3) 결과 메일 ────────────────────────────────────────────
async def maybe_send_result_emails() -> None:
async with SessionLocal() as db:
cutoff = now_utc() - timedelta(minutes=settings.result_email_delay_minutes)
matches = (
await db.execute(
select(Match)
.where(
Match.result_outcome.is_not(None),
Match.finished_at.is_not(None),
Match.results_emailed_at.is_(None),
)
.options(selectinload(Match.predictions))
)
).scalars().all()
for m in matches:
if m.finished_at and ensure_aware(m.finished_at) > cutoff:
continue # 아직 지연시간(3시간) 미경과
conds = [
UserPrediction.match_id == m.match_id,
UserPrediction.notify.is_(True),
UserPrediction.email.is_not(None),
UserPrediction.notified.is_(False),
]
# '맞춘 사람만' 옵션: 승패(outcome) 적중자에게만 발송
if settings.result_email_correct_only:
conds.append(UserPrediction.outcome == m.result_outcome)
picks = (
await db.execute(select(UserPrediction).where(*conds))
).scalars().all()
ai_lines = [(p.model, p.outcome == m.result_outcome) for p in m.predictions]
match_url = f"{settings.public_origin}/match/{m.match_id}"
sent_any = False
for pk in picks:
subject, html, text = build_result_email(
team_a=m.team_a_short,
team_b=m.team_b_short,
result_a=m.result_score_a or 0,
result_b=m.result_score_b or 0,
my_a=pk.score_a,
my_b=pk.score_b,
my_points=pk.points or 0,
ai_lines=ai_lines,
match_url=match_url,
)
try:
await send_email(pk.email, subject, html, text) # type: ignore[arg-type]
pk.notified = True
sent_any = True
except EmailUnavailable as e:
log.warning("result email skipped (%s)", e)
break # SMTP 미설정 — 이 경기는 다음 틱에 재시도
except Exception as e: # noqa: BLE001
log.error("result email failed → %s: %s", pk.email, e)
# 모든 구독자에게 시도 완료(또는 구독자 없음)면 발송완료 표시
still_pending = any(
(not pk.notified) for pk in picks
)
if not still_pending:
m.results_emailed_at = now_utc()
await db.commit()
def build_scheduler() -> AsyncIOScheduler:
sched = AsyncIOScheduler(timezone=settings.timezone)
# 경기 일정: 매일 KST 09:00 외부 크롤링 → 경기·투표시간 갱신
sched.add_job(
sync_schedule_job,
"cron",
hour=settings.schedule_sync_hour,
minute=settings.schedule_sync_minute,
id="schedule_sync",
)
# 투표시간 상태 전이 (open/locked)
sched.add_job(
tick_status,
"interval",
seconds=settings.status_tick_seconds,
id="status_tick",
next_run_time=now_utc(),
)
# 축구 데이터 수집: 매일 KST 00:00 (예측 생성보다 먼저 캐시를 채움)
sched.add_job(
refresh_football_data,
"cron",
hour=settings.football_refresh_hour,
minute=settings.football_refresh_minute,
id="football_refresh",
)
# AI 예측 생성: 매시간 — 킥오프 22h 전에 든 경기를 놓치지 않게.
# only_missing 이라 이미 생성된 (경기×모델)은 API 호출 없이 건너뜀(비용 동일).
sched.add_job(
generate_ai_predictions,
"interval",
hours=1,
id="ai_generate",
next_run_time=now_utc(),
)
# 결과 자동 정산 + 메일: 킥오프+Nh 지난 경기 스코어 수집·채점·발송
sched.add_job(
settle_matches,
"interval",
seconds=settings.settle_tick_seconds,
id="settle",
next_run_time=now_utc(),
)
# 야구(kbo/mlb): 일정·프리뷰 동기화 매일 09:00 + AI 생성 직전 00:00
sched.add_job(
sync_baseball_job, "cron",
hour=settings.schedule_sync_hour, minute=settings.schedule_sync_minute,
id="baseball_sync",
)
sched.add_job(sync_baseball_job, "cron", hour=0, minute=0, id="baseball_sync_night")
# 야구 결과 정산 폴링
sched.add_job(
settle_baseball, "interval",
seconds=settings.settle_tick_seconds, id="baseball_settle",
)
return sched
async def main() -> None:
load_scoring_data() # data/scoring.json → 배점·배제 대상 (채점·결과메일에 사용)
await init_db() # 테이블 보장 (idempotent). 시드는 API 가 담당.
# 기동 시 1회: 일정 동기화(축구+야구) → AI 예측 생성 → 밀린 경기 자동 정산
await sync_schedule_job()
await sync_baseball_job()
await refresh_football_data() # 축구 데이터 캐시 선채움 (키 없으면 no-op)
await generate_ai_predictions()
await settle_matches()
await settle_baseball()
sched = build_scheduler()
sched.start()
log.info(
"scheduling-server started — schedule sync daily %02d:%02d · status every %ss · "
"AI gen hourly (kickoff-%dh) · auto-settle on FINISHED (poll every %ss) · 맞춘사람만=%s",
settings.schedule_sync_hour,
settings.schedule_sync_minute,
settings.status_tick_seconds,
settings.ai_generate_lookahead_hours,
settings.settle_tick_seconds,
settings.result_email_correct_only,
)
# 영구 대기
stop = asyncio.Event()
try:
await stop.wait()
finally:
sched.shutdown()
if __name__ == "__main__":
asyncio.run(main())