diff --git a/db/connection.py b/db/connection.py index 729ca5f..a65bad2 100644 --- a/db/connection.py +++ b/db/connection.py @@ -17,7 +17,11 @@ f"@{os.getenv('DB_HOST')}:{os.getenv('DB_PORT')}/{os.getenv('DB_NAME')}" ) -engine = create_engine(_DATABASE_URL, pool_pre_ping=True) +engine = create_engine( + _DATABASE_URL, + pool_pre_ping=True, + connect_args={"options": "-c statement_timeout=10000 -c lock_timeout=5000"}, +) SessionLocal = sessionmaker(bind=engine, autocommit=False, autoflush=False) diff --git a/db/repository.py b/db/repository.py index d61b52d..e79f6be 100644 --- a/db/repository.py +++ b/db/repository.py @@ -302,44 +302,54 @@ def get_completed_concerts() -> list[dict]: def save_setlists(setlists: list[dict]) -> None: - """수집된 셋리스트를 setlist·setlist_track 테이블에 저장한다. 중복 시 무시.""" - with get_session() as session: - for item in setlists: - row = session.execute( - text(""" - INSERT INTO setlist (concert_id, setlist_fm_id, attribution_url, collected_at) - VALUES (:concert_id, :setlist_fm_id, :attribution_url, NOW()) - ON CONFLICT (setlist_fm_id) DO NOTHING - RETURNING id - """), - { - "concert_id": item["concert_id"], - "setlist_fm_id": item["setlist_fm_id"], - "attribution_url": item.get("attribution_url"), - }, - ).fetchone() - - if not row or not item.get("tracks"): - continue - - setlist_id = row[0] - session.execute( - text(""" - INSERT INTO setlist_track (setlist_id, position, song_name, info) - VALUES (:setlist_id, :position, :song_name, :info) - """), - [ + """수집된 셋리스트를 setlist·setlist_track 테이블에 저장한다. 중복 시 무시. + + 셋리스트별 독립 트랜잭션을 사용해 한 건 실패가 전체에 영향을 주지 않는다. + """ + saved = 0 + for item in setlists: + try: + with get_session() as session: + row = session.execute( + text(""" + INSERT INTO setlist + (concert_id, setlist_fm_id, attribution_url, collected_at) + VALUES (:concert_id, :setlist_fm_id, :attribution_url, NOW()) + ON CONFLICT (setlist_fm_id) DO NOTHING + RETURNING id + """), { - "setlist_id": setlist_id, - "position": track["position"], - "song_name": track["song_name"], - "info": track.get("info"), - } - for track in item["tracks"] - ], + "concert_id": item["concert_id"], + "setlist_fm_id": item["setlist_fm_id"], + "attribution_url": item.get("attribution_url"), + }, + ).fetchone() + + if row and item.get("tracks"): + setlist_id = row[0] + session.execute( + text(""" + INSERT INTO setlist_track (setlist_id, position, song_name, info) + VALUES (:setlist_id, :position, :song_name, :info) + """), + [ + { + "setlist_id": setlist_id, + "position": track["position"], + "song_name": track["song_name"], + "info": track.get("info"), + } + for track in item["tracks"] + ], + ) + saved += 1 + except SQLAlchemyError as e: + logger.error( + "셋리스트 저장 실패 — 건너뜀: setlist_fm_id=%s, 오류=%s", + item.get("setlist_fm_id"), e, ) - logger.info("셋리스트 저장 완료: %d건 처리", len(setlists)) + logger.info("셋리스트 저장 완료: %d / %d건 처리", saved, len(setlists)) def get_all_aliases() -> list[dict]: @@ -619,40 +629,62 @@ def get_unmatched_concerts() -> list[dict]: def save_concert_artists(matches: list[dict]) -> None: - """매칭 결과를 concert_artist 테이블에 저장한다.""" - with get_session() as session: - for match in matches: - session.execute( - text(""" - INSERT INTO concert_artist (concert_id, artist_id) - VALUES (:concert_id, :artist_id) - ON CONFLICT (concert_id, artist_id) DO NOTHING - """), - { - "concert_id": match["concert_id"], - "artist_id": match["artist_id"], - }, + """매칭 결과를 concert_artist 테이블에 저장한다. + + 매칭별 독립 트랜잭션을 사용해 한 건 실패가 전체에 영향을 주지 않는다. + """ + saved = 0 + for match in matches: + try: + with get_session() as session: + session.execute( + text(""" + INSERT INTO concert_artist (concert_id, artist_id) + VALUES (:concert_id, :artist_id) + ON CONFLICT (concert_id, artist_id) DO NOTHING + """), + { + "concert_id": match["concert_id"], + "artist_id": match["artist_id"], + }, + ) + saved += 1 + except SQLAlchemyError as e: + logger.error( + "공연-아티스트 매칭 저장 실패 — 건너뜀: concert_id=%s, artist_id=%s, 오류=%s", + match["concert_id"], match["artist_id"], e, ) - logger.info("공연-아티스트 매칭 저장 완료: %d건 처리", len(matches)) + logger.info("공연-아티스트 매칭 저장 완료: %d / %d건 처리", saved, len(matches)) def save_concert_artist_candidates(matches: list[dict]) -> None: - """자동 매칭 결과를 concert_artist_candidate 테이블에 저장한다. 중복 시 무시.""" - with get_session() as session: - for match in matches: - session.execute( - text(""" - INSERT INTO concert_artist_candidate - (concert_id, artist_id) - VALUES (:concert_id, :artist_id) - ON CONFLICT (concert_id, artist_id) DO NOTHING - """), - { - "concert_id": match["concert_id"], - "artist_id": match["artist_id"], - }, + """자동 매칭 결과를 concert_artist_candidate 테이블에 저장한다. 중복 시 무시. + + 매칭별 독립 트랜잭션을 사용해 한 건 실패가 전체에 영향을 주지 않는다. + """ + saved = 0 + for match in matches: + try: + with get_session() as session: + session.execute( + text(""" + INSERT INTO concert_artist_candidate + (concert_id, artist_id) + VALUES (:concert_id, :artist_id) + ON CONFLICT (concert_id, artist_id) DO NOTHING + """), + { + "concert_id": match["concert_id"], + "artist_id": match["artist_id"], + }, + ) + saved += 1 + except SQLAlchemyError as e: + logger.error( + "공연-아티스트 후보 저장 실패 — 건너뜀: concert_id=%s, artist_id=%s, 오류=%s", + match["concert_id"], match["artist_id"], e, ) - logger.info("공연-아티스트 후보 저장 완료: %d건 처리", len(matches)) + logger.info("공연-아티스트 후보 저장 완료: %d / %d건 처리", saved, len(matches)) @@ -722,31 +754,39 @@ def update_concert_fetch_attempted(concert_id: int) -> None: def update_concert_status(concerts: list[dict]) -> None: - """updatedate 변화 감지 시 status와 kopis_update_date를 갱신한다.""" + """updatedate 변화 감지 시 status와 kopis_update_date를 갱신한다. + + 공연별 독립 트랜잭션을 사용해 한 건 실패가 전체에 영향을 주지 않는다. + """ updated = 0 - with get_session() as session: - for concert in concerts: - row = session.execute( - text("SELECT kopis_update_date FROM concert WHERE kopis_id = :kopis_id"), - {"kopis_id": concert["kopis_id"]}, - ).fetchone() + for concert in concerts: + try: + with get_session() as session: + row = session.execute( + text("SELECT kopis_update_date FROM concert WHERE kopis_id = :kopis_id"), + {"kopis_id": concert["kopis_id"]}, + ).fetchone() - if row is None: - continue + if row is None: + continue - if row[0] != concert["updatedate"]: - session.execute( - text(""" - UPDATE concert - SET status = :status, kopis_update_date = :kopis_update_date - WHERE kopis_id = :kopis_id - """), - { - "status": _KOPIS_STATUS_MAP.get(concert["prfstate"], "PENDING"), - "kopis_update_date": concert["updatedate"], - "kopis_id": concert["kopis_id"], - }, - ) - updated += 1 + if row[0] != concert["updatedate"]: + session.execute( + text(""" + UPDATE concert + SET status = :status, kopis_update_date = :kopis_update_date + WHERE kopis_id = :kopis_id + """), + { + "status": _KOPIS_STATUS_MAP.get(concert["prfstate"], "PENDING"), + "kopis_update_date": concert["updatedate"], + "kopis_id": concert["kopis_id"], + }, + ) + updated += 1 + except SQLAlchemyError as e: + logger.error( + "공연 상태 갱신 실패 — 건너뜀: kopis_id=%s, 오류=%s", concert["kopis_id"], e + ) logger.info("공연 상태 갱신 완료: %d건 변경", updated) diff --git a/scheduler.py b/scheduler.py index b633862..9238e6b 100644 --- a/scheduler.py +++ b/scheduler.py @@ -6,6 +6,7 @@ import threading import time from concurrent.futures import ThreadPoolExecutor +from contextlib import contextmanager from datetime import datetime, timedelta from typing import List, Optional, Set @@ -65,6 +66,16 @@ _scheduler: Optional[BackgroundScheduler] = None # main()에서 데몬 실행 시에만 할당됨 +@contextmanager +def _job_timer(job_name: str): + """잡의 소요시간을 로깅한다. 조기 반환·예외 등 모든 종료 경로를 포함한다.""" + start = time.monotonic() + try: + yield + finally: + logger.info("%s 소요시간: %.1f초", job_name, time.monotonic() - start) + + def _progress_path(job_name: str) -> str: return os.path.join(_CHECKPOINT_DIR, f"{job_name}.json") @@ -356,44 +367,45 @@ def run_ja_romanize_collect() -> None: locale='ko' alias가 이미 존재하는 아티스트는 건너뛴다. 외부 API 호출 없이 DB에 이미 저장된 sort_name만으로 동작한다. """ - logger.info("=== 로마자→한글 alias 변환 잡 시작 ===") - artists = get_artists_without_ko_alias() - logger.info("한국어 alias 미수집 아티스트: %d건", len(artists)) - if not artists: - logger.info("=== 로마자→한글 alias 변환 잡 완료 (대상 없음) ===") - return - aliases = ja_romanize.collect_ko_aliases(artists) - if aliases: - save_aliases(aliases) - logger.info("=== 로마자→한글 alias 변환 잡 완료 ===") + with _job_timer("로마자→한글 alias 변환 잡"): + logger.info("=== 로마자→한글 alias 변환 잡 시작 ===") + artists = get_artists_without_ko_alias() + logger.info("한국어 alias 미수집 아티스트: %d건", len(artists)) + if not artists: + logger.info("=== 로마자→한글 alias 변환 잡 완료 (대상 없음) ===") + return + aliases = ja_romanize.collect_ko_aliases(artists) + if aliases: + save_aliases(aliases) + logger.info("=== 로마자→한글 alias 변환 잡 완료 ===") def run_concert_status_update() -> None: """활성 공연 상태 갱신 (매일). DB의 진행 중 공연을 개별 API로 최신 상태로 갱신한다.""" - logger.info("=== 공연 상태 갱신 잡 시작 ===") - - active = get_active_concerts() - if active: - logger.info("활성 공연 %d건 상태 갱신 시작", len(active)) - with ThreadPoolExecutor(max_workers=8) as pool: - futures = { - pool.submit(kopis.collect_by_id, c["kopis_id"]): c["kopis_id"] - for c in active - } - fetched = [] - for future, kopis_id in futures.items(): - try: - result = future.result() - if result is not None: - fetched.append(result) - except Exception as e: - logger.warning("상태 갱신 API 실패 kopis_id=%s: %s", kopis_id, e) - if fetched: - update_concert_status(fetched) - - update_artist_is_coming() - logger.info("=== 공연 상태 갱신 잡 완료 ===") + with _job_timer("공연 상태 갱신 잡"): + logger.info("=== 공연 상태 갱신 잡 시작 ===") + + active = get_active_concerts() + if active: + logger.info("활성 공연 %d건 상태 갱신 시작", len(active)) + with ThreadPoolExecutor(max_workers=8) as pool: + futures = { + pool.submit(kopis.collect_by_id, c["kopis_id"]): c["kopis_id"] + for c in active + } + fetched = [] + for future, kopis_id in futures.items(): + try: + result = future.result() + if result is not None: + fetched.append(result) + except Exception as e: + logger.warning("상태 갱신 API 실패 kopis_id=%s: %s", kopis_id, e) + if fetched: + update_concert_status(fetched) + + logger.info("=== 공연 상태 갱신 잡 완료 ===") def run_new_concert_collect(stdate: Optional[str] = None, use_prfstate: bool = False) -> None: @@ -407,53 +419,56 @@ def run_new_concert_collect(stdate: Optional[str] = None, use_prfstate: bool = F 어드민 검수 큐로 보낸다. use_prfstate=True(초기 수집 전용)이면 KOPIS prfstate를 그대로 반영해 검수 없이 실제 상태로 저장한다. """ - logger.info("=== 신규 공연 탐지 잡 시작 ===") + with _job_timer("신규 공연 탐지 잡"): + logger.info("=== 신규 공연 탐지 잡 시작 ===") - last_date = _load_last_collect_date() - if last_date: - logger.info("증분 스캔 — afterdate=%s", last_date) - else: - logger.info("초기 전체 스캔 (체크포인트 없음, stdate=%s)", stdate or "20250101") - concerts = kopis.collect(stdate=stdate, afterdate=last_date) - existing_ids = get_existing_kopis_ids() - aliases = get_all_aliases() - new_concerts = [ - c for c in concerts - if c["kopis_id"] not in existing_ids and has_match(c, aliases) - ] - if new_concerts: - logger.info("신규 공연 %d건 저장 시작", len(new_concerts)) - save_concerts(new_concerts, use_prfstate=use_prfstate) - unmatched = get_unmatched_concerts() - all_matches: list[dict] = [] - for concert in unmatched: - matches, _ = match_concert(concert, aliases) - all_matches.extend(matches) - if all_matches: - save_concert_artist_candidates(all_matches) - - new_kopis_titles = {c["kopis_id"]: c["prfnm"] for c in new_concerts} - concert_ids_by_kopis_id = get_concert_ids_by_kopis_ids(list(new_kopis_titles)) - title_by_concert_id = { - concert_id: new_kopis_titles[kopis_id] - for kopis_id, concert_id in concert_ids_by_kopis_id.items() - } + last_date = _load_last_collect_date() + if last_date: + logger.info("증분 스캔 — afterdate=%s", last_date) + else: + logger.info("초기 전체 스캔 (체크포인트 없음, stdate=%s)", stdate or "20250101") + concerts = kopis.collect(stdate=stdate, afterdate=last_date) + existing_ids = get_existing_kopis_ids() + aliases = get_all_aliases() + new_concerts = [ + c for c in concerts + if c["kopis_id"] not in existing_ids and has_match(c, aliases) + ] + if new_concerts: + logger.info("신규 공연 %d건 저장 시작", len(new_concerts)) + save_concerts(new_concerts, use_prfstate=use_prfstate) + unmatched = get_unmatched_concerts() + all_matches: list[dict] = [] + for concert in unmatched: + matches, _ = match_concert(concert, aliases) + all_matches.extend(matches) + if all_matches: + save_concert_artist_candidates(all_matches) + + new_kopis_titles = {c["kopis_id"]: c["prfnm"] for c in new_concerts} + concert_ids_by_kopis_id = get_concert_ids_by_kopis_ids(list(new_kopis_titles)) + title_by_concert_id = { + concert_id: new_kopis_titles[kopis_id] + for kopis_id, concert_id in concert_ids_by_kopis_id.items() + } - matches_by_concert: dict = {} - for match in all_matches: - if match["concert_id"] in title_by_concert_id: - matches_by_concert.setdefault(match["concert_id"], []).append(match["artist_id"]) + matches_by_concert: dict = {} + for match in all_matches: + if match["concert_id"] in title_by_concert_id: + matches_by_concert.setdefault( + match["concert_id"], [] + ).append(match["artist_id"]) - if matches_by_concert: - all_artist_ids = {aid for ids in matches_by_concert.values() for aid in ids} - artist_names = get_artist_names_by_ids(list(all_artist_ids)) - for concert_id, artist_ids in matches_by_concert.items(): - names = [artist_names[aid] for aid in artist_ids if aid in artist_names] - notify_new_concert(title_by_concert_id[concert_id], names) + if matches_by_concert: + all_artist_ids = {aid for ids in matches_by_concert.values() for aid in ids} + artist_names = get_artist_names_by_ids(list(all_artist_ids)) + for concert_id, artist_ids in matches_by_concert.items(): + names = [artist_names[aid] for aid in artist_ids if aid in artist_names] + notify_new_concert(title_by_concert_id[concert_id], names) - update_artist_is_coming() - _save_last_collect_date(datetime.now().strftime("%Y%m%d")) - logger.info("=== 신규 공연 탐지 잡 완료 ===") + update_artist_is_coming() + _save_last_collect_date(datetime.now().strftime("%Y%m%d")) + logger.info("=== 신규 공연 탐지 잡 완료 ===") @@ -464,81 +479,85 @@ def run_artist_image_update() -> None: 마지막 실패일을 기록해두고, _IMAGE_RETRY_DAYS가 지나기 전까지는 재시도 대상에서 제외한다 — 매주 동일한 잔여 아티스트로 Spotify API를 반복 호출하지 않기 위함. """ - logger.info("=== 아티스트 이미지 수집 잡 시작 ===") - if _is_banned(): - logger.warning("Spotify 429 밴 유효 — 이미지 수집 건너뜀") - return - artists = get_artists_without_image() - failed = _load_failed_image_artists() - cutoff = (datetime.now() - timedelta(days=_IMAGE_RETRY_DAYS)).strftime("%Y%m%d") - targets = [a for a in artists if failed.get(a["id"], "") < cutoff] - skipped = len(artists) - len(targets) - if skipped: - logger.info("최근 %d일 내 실패 기록 — 재시도 보류: %d건", _IMAGE_RETRY_DAYS, skipped) - logger.info("이미지 수집 시도 대상: %d건", len(targets)) - - today = datetime.now().strftime("%Y%m%d") - for a in targets: - try: - image_url, spotify_id = artist_image.collect_artist_image( - a["mbid"], - spotify_url=a.get("spotify_url"), - name=a.get("name"), - ) - except SpotifyRateLimitError as e: - logger.error("Spotify 429 — 이미지 수집 중단 (artist_id=%d)", a["id"]) - banned_until = _save_ban(e.retry_after) - _reschedule_on_ban_lift(run_artist_image_update, banned_until) - break - except Exception as e: - logger.warning("아티스트 이미지 수집 실패 mbid=%s: %s", a["mbid"], e) - continue - if image_url: - update_artist_image(a["id"], image_url) - failed.pop(a["id"], None) - else: - failed[a["id"]] = today - if not a.get("spotify_url") and spotify_id: - upsert_artist_url(a["id"], "Spotify", artist_image.spotify_artist_url(spotify_id)) - _save_failed_image_artists(failed) - logger.info("=== 아티스트 이미지 수집 잡 완료 ===") + with _job_timer("아티스트 이미지 수집 잡"): + logger.info("=== 아티스트 이미지 수집 잡 시작 ===") + if _is_banned(): + logger.warning("Spotify 429 밴 유효 — 이미지 수집 건너뜀") + return + artists = get_artists_without_image() + failed = _load_failed_image_artists() + cutoff = (datetime.now() - timedelta(days=_IMAGE_RETRY_DAYS)).strftime("%Y%m%d") + targets = [a for a in artists if failed.get(a["id"], "") < cutoff] + skipped = len(artists) - len(targets) + if skipped: + logger.info("최근 %d일 내 실패 기록 — 재시도 보류: %d건", _IMAGE_RETRY_DAYS, skipped) + logger.info("이미지 수집 시도 대상: %d건", len(targets)) + + today = datetime.now().strftime("%Y%m%d") + for a in targets: + try: + image_url, spotify_id = artist_image.collect_artist_image( + a["mbid"], + spotify_url=a.get("spotify_url"), + name=a.get("name"), + ) + except SpotifyRateLimitError as e: + logger.error("Spotify 429 — 이미지 수집 중단 (artist_id=%d)", a["id"]) + banned_until = _save_ban(e.retry_after) + _reschedule_on_ban_lift(run_artist_image_update, banned_until) + break + except Exception as e: + logger.warning("아티스트 이미지 수집 실패 mbid=%s: %s", a["mbid"], e) + continue + if image_url: + update_artist_image(a["id"], image_url) + failed.pop(a["id"], None) + else: + failed[a["id"]] = today + if not a.get("spotify_url") and spotify_id: + upsert_artist_url( + a["id"], "Spotify", artist_image.spotify_artist_url(spotify_id) + ) + _save_failed_image_artists(failed) + logger.info("=== 아티스트 이미지 수집 잡 완료 ===") def run_release_update() -> None: """릴리즈 갱신 (매일). 내한 공연 매칭 아티스트만 대상.""" - logger.info("=== 릴리즈 갱신 잡 시작 ===") - if _is_banned(): - logger.warning("Spotify 429 밴 유효 — 릴리즈 갱신 건너뜀") - return - completed_ids = _load_progress("release_update") - for a in get_matched_artists_with_spotify(): - artist_id = a["artist_id"] - if artist_id in completed_ids: - logger.info("체크포인트 — 완료 아티스트 건너뜀: artist_id=%d", artist_id) - continue - try: - spotify_id = _extract_spotify_id(a["spotify_url"]) - cached_total = get_spotify_album_total(artist_id) - existing_ids = get_existing_release_spotify_ids(artist_id) - raw, spotify_total = release.collect_releases( - spotify_id, skip_spotify_ids=existing_ids, cached_total=cached_total - ) - if spotify_total > 0 and spotify_total != cached_total: - update_spotify_album_total(artist_id, spotify_total) - if raw: - save_releases(artist_id, _sort_releases(raw)) - completed_ids.add(artist_id) - except SpotifyRateLimitError as e: - logger.error("Spotify 429 — 릴리즈 갱신 중단 (artist_id=%d)", artist_id) - banned_until = _save_ban(e.retry_after) - _save_progress("release_update", completed_ids) - _reschedule_on_ban_lift(run_release_update, banned_until) - break - except Exception as e: - logger.error("릴리즈 수집 실패 — artist_id=%d: %s", artist_id, e) - else: - _clear_progress("release_update") - logger.info("=== 릴리즈 갱신 잡 완료 ===") + with _job_timer("릴리즈 갱신 잡"): + logger.info("=== 릴리즈 갱신 잡 시작 ===") + if _is_banned(): + logger.warning("Spotify 429 밴 유효 — 릴리즈 갱신 건너뜀") + return + completed_ids = _load_progress("release_update") + for a in get_matched_artists_with_spotify(): + artist_id = a["artist_id"] + if artist_id in completed_ids: + logger.info("체크포인트 — 완료 아티스트 건너뜀: artist_id=%d", artist_id) + continue + try: + spotify_id = _extract_spotify_id(a["spotify_url"]) + cached_total = get_spotify_album_total(artist_id) + existing_ids = get_existing_release_spotify_ids(artist_id) + raw, spotify_total = release.collect_releases( + spotify_id, skip_spotify_ids=existing_ids, cached_total=cached_total + ) + if spotify_total > 0 and spotify_total != cached_total: + update_spotify_album_total(artist_id, spotify_total) + if raw: + save_releases(artist_id, _sort_releases(raw)) + completed_ids.add(artist_id) + except SpotifyRateLimitError as e: + logger.error("Spotify 429 — 릴리즈 갱신 중단 (artist_id=%d)", artist_id) + banned_until = _save_ban(e.retry_after) + _save_progress("release_update", completed_ids) + _reschedule_on_ban_lift(run_release_update, banned_until) + break + except Exception as e: + logger.error("릴리즈 수집 실패 — artist_id=%d: %s", artist_id, e) + else: + _clear_progress("release_update") + logger.info("=== 릴리즈 갱신 잡 완료 ===") def run_missing_release_update() -> None: @@ -648,10 +667,11 @@ def register_artist_by_mbid(mbid: str) -> dict: def run_setlist_collect() -> None: """setlist.fm 수집 (매일).""" - logger.info("=== setlist 수집 잡 시작 ===") - setlists = setlist.collect() - save_setlists(setlists) - logger.info("=== setlist 수집 잡 완료 ===") + with _job_timer("setlist 수집 잡"): + logger.info("=== setlist 수집 잡 시작 ===") + setlists = setlist.collect() + save_setlists(setlists) + logger.info("=== setlist 수집 잡 완료 ===") def collect_and_save_concert(kopis_id: str) -> dict: diff --git a/tests/test_kopis.py b/tests/test_kopis.py index 25f394b..7e84b6c 100644 --- a/tests/test_kopis.py +++ b/tests/test_kopis.py @@ -5,6 +5,7 @@ import pytest import requests +from sqlalchemy.exc import SQLAlchemyError from collectors.kopis import ( _fetch_and_merge, @@ -869,6 +870,27 @@ def test_updates_prfstate_and_updatedate_when_changed(self): assert update_params["status"] == "ENDED" assert update_params["kopis_update_date"] == "2024-01-20" + def test_one_failure_does_not_stop_others(self): + """한 건이 SQLAlchemyError로 실패해도 나머지 건은 계속 처리되어야 한다.""" + mock_session = MagicMock() + exec_result_1 = MagicMock() + exec_result_1.fetchone.return_value = ("2024-01-01",) + exec_result_3 = MagicMock() + exec_result_3.fetchone.return_value = ("2024-01-03",) + mock_session.execute.side_effect = [exec_result_1, SQLAlchemyError("boom"), exec_result_3] + concerts = [ + {"kopis_id": "PF001", "prfstate": "공연예정", "updatedate": "2024-01-01"}, + {"kopis_id": "PF002", "prfstate": "공연예정", "updatedate": "2024-01-02"}, + {"kopis_id": "PF003", "prfstate": "공연예정", "updatedate": "2024-01-03"}, + ] + + with patch("db.repository.get_session") as mock_get_session: + mock_get_session.return_value.__enter__ = MagicMock(return_value=mock_session) + mock_get_session.return_value.__exit__ = MagicMock(return_value=False) + update_concert_status(concerts) + + assert mock_session.execute.call_count == 3 + # ─── TestKopisHttpErrors ────────────────────────────────────────────────────── diff --git a/tests/test_repository.py b/tests/test_repository.py index 7975535..1b72141 100644 --- a/tests/test_repository.py +++ b/tests/test_repository.py @@ -1,10 +1,13 @@ from unittest.mock import MagicMock, patch +from sqlalchemy.exc import SQLAlchemyError + from db.repository import ( get_active_concerts, get_artist_names_by_ids, get_concert_ids_by_kopis_ids, save_artists, + save_concert_artist_candidates, save_concert_artists, save_setlists, update_artist_is_coming, @@ -414,6 +417,22 @@ def test_skips_tracks_when_conflict_returns_no_row(self): sqls = [str(c.args[0]) for c in mock_session.execute.call_args_list] assert not any("setlist_track" in s for s in sqls) + def test_one_failure_does_not_stop_others(self): + """한 건이 SQLAlchemyError로 실패해도 나머지 건은 계속 저장 시도되어야 한다.""" + mock_session = MagicMock() + # tracks=[]인 항목은 execute().fetchone() 반환값이 사용되지 않으므로 + # 성공 케이스는 임의의 MagicMock, 실패 케이스만 예외로 지정한다. + mock_session.execute.side_effect = [MagicMock(), SQLAlchemyError("boom"), MagicMock()] + setlists = [ + {"concert_id": 1, "setlist_fm_id": "abc1", "tracks": []}, + {"concert_id": 2, "setlist_fm_id": "abc2", "tracks": []}, + {"concert_id": 3, "setlist_fm_id": "abc3", "tracks": []}, + ] + + self._run(setlists, mock_session) + + assert mock_session.execute.call_count == 3 + class TestSaveConcertArtists: def _run(self, matches, mock_session): @@ -472,6 +491,50 @@ def test_handles_empty_list(self): self._run([], mock_session) mock_session.execute.assert_not_called() + def test_one_failure_does_not_stop_others(self): + """한 건이 SQLAlchemyError로 실패해도 나머지 건은 계속 저장 시도되어야 한다.""" + mock_session = MagicMock() + mock_session.execute.side_effect = [None, SQLAlchemyError("boom"), None] + matches = [ + {"concert_id": 1, "artist_id": 10}, + {"concert_id": 2, "artist_id": 20}, + {"concert_id": 3, "artist_id": 30}, + ] + + self._run(matches, mock_session) + + assert mock_session.execute.call_count == 3 + + +class TestSaveConcertArtistCandidates: + def _run(self, matches, mock_session): + with patch("db.repository.get_session") as mock_get_session: + mock_get_session.return_value.__enter__ = MagicMock(return_value=mock_session) + mock_get_session.return_value.__exit__ = MagicMock(return_value=False) + save_concert_artist_candidates(matches) + + def test_inserts_into_concert_artist_candidate(self): + """매칭 결과가 concert_artist_candidate 테이블에 INSERT되어야 한다.""" + mock_session = MagicMock() + self._run([{"concert_id": 1, "artist_id": 10}], mock_session) + + insert_sql = str(mock_session.execute.call_args_list[0].args[0]) + assert "INSERT INTO concert_artist_candidate" in insert_sql + + def test_one_failure_does_not_stop_others(self): + """한 건이 SQLAlchemyError로 실패해도 나머지 건은 계속 저장 시도되어야 한다.""" + mock_session = MagicMock() + mock_session.execute.side_effect = [SQLAlchemyError("boom"), None, None] + matches = [ + {"concert_id": 1, "artist_id": 10}, + {"concert_id": 2, "artist_id": 20}, + {"concert_id": 3, "artist_id": 30}, + ] + + self._run(matches, mock_session) + + assert mock_session.execute.call_count == 3 + class TestUpdateArtistIsComing: def _make_session_mock(self, rowcount=0): diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index e52b3b1..93de9a7 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -1,5 +1,6 @@ """scheduler.py 단위 테스트.""" import json +import logging import os from datetime import datetime, timedelta from unittest.mock import patch @@ -10,6 +11,7 @@ _build_scheduler, _clear_progress, _is_banned, + _job_timer, _load_progress, _progress_path, _save_ban, @@ -29,6 +31,38 @@ ) +class TestJobTimer: + # ── 배치 잡 소요시간 로깅 ──────────────────────────────────────────────── + + def test_logs_duration_on_normal_completion(self, caplog): + """정상 종료 시 소요시간이 로깅되어야 한다.""" + with caplog.at_level(logging.INFO, logger="scheduler"): + with _job_timer("테스트 잡"): + pass + + assert any("테스트 잡 소요시간" in r.message for r in caplog.records) + + def test_logs_duration_on_early_return(self, caplog): + """with 블록 중간에 return으로 빠져나가도 소요시간이 로깅되어야 한다.""" + def _early_return(): + with _job_timer("테스트 잡"): + return + + with caplog.at_level(logging.INFO, logger="scheduler"): + _early_return() + + assert any("테스트 잡 소요시간" in r.message for r in caplog.records) + + def test_logs_duration_even_on_exception(self, caplog): + """예외가 발생해도 소요시간이 로깅되고, 예외는 그대로 전파되어야 한다.""" + with caplog.at_level(logging.INFO, logger="scheduler"): + with pytest.raises(ValueError): + with _job_timer("테스트 잡"): + raise ValueError("boom") + + assert any("테스트 잡 소요시간" in r.message for r in caplog.records) + + class TestRunInitialCollect: @pytest.fixture(autouse=True) def mock_sub_jobs(self): @@ -159,7 +193,6 @@ def test_calls_collect_by_id_for_each_active_concert(self): patch("scheduler.get_active_concerts", return_value=active), patch("scheduler.kopis.collect_by_id", return_value=fetched) as mock_by_id, patch("scheduler.update_concert_status"), - patch("scheduler.update_artist_is_coming"), ): run_concert_status_update() @@ -175,7 +208,6 @@ def test_update_concert_status_called_with_fetched_results(self): patch("scheduler.get_active_concerts", return_value=active), patch("scheduler.kopis.collect_by_id", return_value=fetched), patch("scheduler.update_concert_status") as mock_update, - patch("scheduler.update_artist_is_coming"), ): run_concert_status_update() @@ -187,7 +219,6 @@ def test_skips_status_update_when_no_active_concerts(self): patch("scheduler.get_active_concerts", return_value=[]), patch("scheduler.kopis.collect_by_id") as mock_by_id, patch("scheduler.update_concert_status") as mock_update, - patch("scheduler.update_artist_is_coming"), ): run_concert_status_update() @@ -201,21 +232,20 @@ def test_skips_none_results_from_collect_by_id(self): patch("scheduler.get_active_concerts", return_value=active), patch("scheduler.kopis.collect_by_id", return_value=None), patch("scheduler.update_concert_status") as mock_update, - patch("scheduler.update_artist_is_coming"), ): run_concert_status_update() mock_update.assert_not_called() - def test_updates_is_coming_after_status_update(self): - """상태 갱신 후 update_artist_is_coming이 인자 없이 호출되어야 한다.""" + def test_does_not_update_is_coming(self): + """is_coming 갱신은 run_new_concert_collect로 일원화됐으므로 이 잡에서는 호출 안 된다.""" with ( patch("scheduler.get_active_concerts", return_value=[]), patch("scheduler.update_artist_is_coming") as mock_update, ): run_concert_status_update() - mock_update.assert_called_once_with() + mock_update.assert_not_called() class TestRunNewConcertCollect: