크롤·IP 로테이션을 검수하며 찾은 결함을 순서대로 고쳤다. 의심 지점은 모두 실제 코드 경로로 재현해 확인했다(브라우저·네트워크만 mock, 프록시·DB 장부는 실물). **① 요청 예산이 사실상 발화하지 않았다 (핵심)** 유휴 정리(close_if_idle, 120s)는 브라우저만 닫고 임대는 두는데, 재기동 때마다 _ip_requests 를 0 으로 되돌렸다. 게다가 ensure_port 의 renew 가 임대 만료를 계속 뒤로 민다 — 검색이 유휴 임계보다 뜸하고 임대(10분)보다 잦으면 **한 IP 에 영원히 고정**된다(실측: 6회 검색이 전부 같은 포트·ip_req#1). 연속 검색에서는 정상 동작해 부하 테스트로는 안 잡히고, 수동 트리거처럼 드문드문한 실사용 패턴에서만 깨진다. 파급이 하나 더 있다 — bot_detection.ip_request_no 가 항상 1 로 찍혀, operations.md 가 명시한 '1 위주면 IP 평판 / 2 이상이면 예산 하향' 진단이 통째로 무너진다. 과거 "전량 ip_req#1 이라 IP 평판 문제" 결론은 이 착시일 수 있다(문서에 경고 추가). → IP 세션 상태를 브라우저 수명과 분리. **포트가 실제로 바뀔 때만** 리셋한다(_begin_ip_session). 세션 종료 기록도 포트 변경·최종 close 시점으로 옮겼다(idle 사유 소멸). **② 환경 차단이면 회복 못 하는데 풀을 계속 태웠다** 쿠팡은 fatal 마커가 없어 컨테이너 차단 같은 '회전 무효' 상황을 구분 못 했다. 실측으로 웜업 6포트 + 잡 1건당 6포트를 30분 쿨다운에 묶어 **잡 16건이면 100포트 고갈**. 실제 장부에도 9분간 11포트 연속 소각 이력이 남아 있다(gate 사용 21 / 소각 14). → 서킷브레이커: **서로 다른 IP 가 연속 3개 모두 첫 요청부터** 막히면 IP 문제가 아니라고 판정, 태우기를 멈추고 fatal 로 알린다(env_block 마커 → 기존 fatal_block 알림이 집계). 같은 IP 반복 차단·뒤쪽 요청 차단은 세지 않는다. 성공 1회로 자동 해제(타이머 불필요). 결과: 전면 차단 시 소각이 판정 근거 2개에서 멈춘다(웜업 6→0, 잡 6→0). **③ 종료가 임대를 반납하지 않았다** close() 후에도 leased_until(최대 10분)까지 그 IP 를 아무도 못 썼다 — 재시작이 잦을수록 가용 풀이 줄었다. DecodoProxy.release() 추가, close() 에서만 호출(유휴 정리는 웜 쿠키·예산 유지를 위해 그대로 둔다). **④ 시간창 재기동이 IP 를 안 바꿨다** — 로그만 'IP 회전'이었고 renew 로 같은 포트를 붙잡았다. sticky 수명이 끝나면 같은 포트라도 IP 가 바뀌므로 명시적으로 놓아준다. **⑤ '검색결과 없음'을 차단으로 오인해 IP 를 태울 수 있었다** 네이버 무결과 페이지 크기는 실측된 적이 없는데 short_html 폴백이 이를 차단으로 본다. 확신도로 대응을 갈랐다 — 알려진 마커만 태우고/서킷브레이커에 세고, 미지의 짧은 HTML 은 회전·재시도까지만. 판단 근거는 bot_detection 에 계속 쌓이므로 나중에 임계를 실측할 수 있다. **⑥** available_ports() 가 장부 모드에서 늘 최대값을 반환하는 점을 문서화(관측 경로는 미사용). 세션 마감을 멱등하게 만들어 close() 중복 호출 시 이중 기록 방지. 테스트 14건 추가(전체 251 passed). mock 하니스도 실물을 타도록 고쳤다 — 회전 시 포트가 실제로 바뀌고, 재기동 판단·IP 세션 경계는 실제 코드를 그대로 쓴다(고정 포트 mock 은 이 버그를 못 봤다). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
395 lines
25 KiB
Python
395 lines
25 KiB
Python
# LPS 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일).
|
||
# python worker_main.py
|
||
# WORKER_CONCURRENCY=3 python worker_main.py # 상품 3개 동시 검색(권장 2~3, 로컬)
|
||
#
|
||
# 브라우저 어댑터는 컨텍스트당 직렬(lock)이라, 진짜 병렬을 위해 **워커마다 자기 브라우저 세트**를 준다:
|
||
# 프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP). 동시성 N → 최대 4×N Chrome.
|
||
|
||
import asyncio
|
||
import os
|
||
import signal
|
||
import time
|
||
|
||
from common.alerts import AlertManager
|
||
from common.database.db_session_manager import DB_SESSION_MNG
|
||
from common.logger import LOG
|
||
from config.server_configs import web_server_config, openai_config, decodo_config, worker_config, alert_config
|
||
from crud.job_crud import JobQueue
|
||
from crud.negative_cache import NegativeCache
|
||
from crud.bot_detection import BotDetectionLog
|
||
from crud.ip_session import IpSessionLog
|
||
from crud.port_lease import PortLeaseStore
|
||
from crud.price_history import PriceHistory
|
||
from services.search.browser_base import ENV_BLOCK_MARKER
|
||
from services.search.profile_slot import claim_profile_slot
|
||
from services.search.proxy import DecodoProxy
|
||
from services.search.coupang.adapter import CoupangAdapter
|
||
from services.search.naver_shop.adapter import NaverShopAdapter
|
||
from services.search.esm.adapter import EsmAdapter
|
||
from services.search.st11.adapter import ElevenStAdapter
|
||
from services.ai.similarity import SimilarityJudge
|
||
from services.ai.keyword import KeywordGenerator
|
||
from worker.handlers import build_search_handler
|
||
from worker.notify import JobListener
|
||
from worker.runner import Worker, run_reaper
|
||
|
||
LOG.SetPrefix(f"{web_server_config.server_name}-worker")
|
||
|
||
# 오픈마켓 폴백(G마켓·옥션·11번가)은 **기본 비활성** — 2026-07-10 협의 결정.
|
||
# 실측상 크롤 몰이 최종 최저가를 바꾼 적이 없고(0회), 검색당 최대 15s + 프록시 대역폭의
|
||
# 대부분을 차지해 로직에서 제외했다(코드·테스트는 유지, 핸들러는 빈 폴백을 정상 처리).
|
||
# 재가동: [WorkerConfig].fallbacks = ["gmarket", "auction", "st11"] (일부만도 가능) —
|
||
# 켜기 전 라이브 스모크로 셀렉터 드리프트 점검. 배경은 docs/decision-openmarket-crawler.md.
|
||
_FALLBACK_SOURCES = ("gmarket", "auction", "st11")
|
||
|
||
|
||
def _enabled_fallbacks() -> list[str]:
|
||
names = [s.strip() for s in worker_config.fallbacks if s.strip()]
|
||
unknown = [n for n in names if n not in _FALLBACK_SOURCES]
|
||
if unknown:
|
||
LOG.w(f"[WorkerConfig].fallbacks 무시된 값: {unknown} (가능: {list(_FALLBACK_SOURCES)})")
|
||
return [n for n in names if n in _FALLBACK_SOURCES]
|
||
|
||
|
||
def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history, store=None):
|
||
"""워커 1개의 자립 세트(브라우저 어댑터·AI·핸들러)를 만든다.
|
||
프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP)로 진짜 병렬을 보장한다."""
|
||
# 워커별 프록시(다른 포트=다른 IP). 100포트를 워커 수로 균등 분할해 시작점을 벌린다.
|
||
# 두 프록시는 **게이트웨이가 다르다** — 쿠팡=국가 무지정, 네이버=한국 전용.
|
||
# 해외 residential IP 로는 네이버 msearch 가 즉시 하드차단된다(2026-08-05 실측:
|
||
# 같은 포트 10091 에서 gate=차단 2,641B / kr=정상 14건). kr_host 가 비면 gate 로
|
||
# 떨어지고 그때는 네이버가 막힌다(로그의 [naver][BOT-DETECTED] 로 드러남).
|
||
# 포트는 **DB 장부(proxy_port)** 가 배타 임대해 준다 — 워커끼리, 두 소스끼리, 그리고
|
||
# **여러 프로세스끼리** 같은 IP 를 동시에 쓰거나 태운 IP 를 곧바로 재사용하는 걸 막는다.
|
||
# owner 에 PID 를 넣어 프로세스가 늘어도 임대자가 구분된다(죽으면 leased_until 로 자동 회수).
|
||
tag = f"p{os.getpid()}-w{i}"
|
||
proxy = DecodoProxy(store=store, owner=f"coupang-{tag}")
|
||
kr_proxy = DecodoProxy(host=decodo_config.kr_host or None, store=store, owner=f"naver-{tag}")
|
||
# 예전엔 seed_offset 으로 워커별 시작 포트를 벌렸지만, 이제 DB 가 LRU(가장 오래 안 쓴 IP)로
|
||
# 배정하므로 프로세스·워커가 몇 개든 알아서 갈린다 — 오프셋 계산이 필요 없다.
|
||
bot_log = BotDetectionLog()
|
||
ip_log = IpSessionLog()
|
||
def _pf(source):
|
||
"""이 워커가 배타적으로 쓸 Chrome 프로필 경로.
|
||
|
||
Chrome 은 프로필당 인스턴스 1개만 허용한다(ProcessSingleton). 예전엔 워커 인덱스로만
|
||
갈랐는데, 그러면 **프로세스가 여러 개일 때 같은 경로를 잡아** 두 번째 프로세스의 검색이
|
||
전부 죽는다(실측). 파일 락으로 슬롯을 선점해 프로세스가 몇 개든 겹치지 않게 한다.
|
||
[WorkerConfig].profile_dir 를 영속 볼륨으로 두면 슬롯 재사용으로 웜 쿠키가 유지된다.
|
||
"""
|
||
return claim_profile_slot(worker_config.profile_dir, source, worker_index=i)
|
||
|
||
adapters = {
|
||
"coupang": CoupangAdapter(headless=False, user_data_dir=_pf("coupang"), proxy=proxy,
|
||
on_detect=bot_log.record, on_session_end=ip_log.record),
|
||
# 네이버는 오픈API(shop.json)가 2026-07-31 종료돼 404 SE05 만 돌아온다 → 모바일 크롤로 전환.
|
||
# 쿠팡과 같은 프록시를 물린다: 단일 IP 로 연달아 두드리면 WTM 캡차로 넘어간다(실측 —
|
||
# 스파이크에서 십수 회 만에 홈 IP 가 탔다). 회전·예산·차단기록은 베이스가 처리한다.
|
||
"naver": NaverShopAdapter(headless=False, user_data_dir=_pf("naver"), proxy=kr_proxy,
|
||
ip_request_budget=decodo_config.naver_ip_request_budget,
|
||
on_detect=bot_log.record, on_session_end=ip_log.record),
|
||
}
|
||
# 폴백은 기본 비활성(LPS_FALLBACKS 로 켬 — 상단 주석 참고). 켤 땐 데드라인이 상한이라
|
||
# 봇감지 재시도(챌린지 대기 2배)를 끈다(max_block_retries=0) — 빠르게 포기·스킵.
|
||
fallback_adapters = {}
|
||
for name in _enabled_fallbacks():
|
||
if name == "st11":
|
||
fallback_adapters[name] = ElevenStAdapter(headless=False, user_data_dir=_pf("st11"), proxy=proxy,
|
||
on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0)
|
||
else:
|
||
fallback_adapters[name] = EsmAdapter(name, headless=False, user_data_dir=_pf(name), proxy=proxy,
|
||
on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0)
|
||
# AI 도 워커별 인스턴스 — 공유 상태(last_usage) 경합 원천 제거
|
||
judge = SimilarityJudge() if has_openai else None
|
||
keyword_gen = KeywordGenerator() if has_openai else None
|
||
handler = build_search_handler(
|
||
adapters, judge=judge, keyword_gen=keyword_gen,
|
||
neg_cache=neg_cache, history=history,
|
||
fallback_adapters=fallback_adapters,
|
||
ai_model=openai_config.model,
|
||
proxy_cost_per_gb=decodo_config.cost_per_gb,
|
||
)
|
||
return handler, list(adapters.values()) + list(fallback_adapters.values())
|
||
|
||
|
||
# 웜업 대상 = 브라우저로 긁는 소스 전부. 챌린지 쿠키 선점이 목적이지만, 동시에
|
||
# **기동 직후 크롤 가능 여부를 확인하는 프리플라이트**이기도 하다 — 여기서 실패하면
|
||
# 그 소스는 이 환경에서 아예 못 긁는다는 뜻이라, 잡이 쌓여 DEAD 될 때까지 기다리지 않고 바로 알린다.
|
||
# (실측 배경: 컨테이너 워커는 쿠팡·네이버가 모두 차단되는데 프로세스는 healthy 라 조용히 0건이 된다)
|
||
_WARMUP_SOURCES = ("gmarket", "auction", "coupang", "naver")
|
||
|
||
|
||
async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float = 60.0, alerts=None):
|
||
"""워커의 크롤 소스를 미리 한 번 긁어 (1) 챌린지 쿠키 확보 (2) 크롤 가능 여부 확인.
|
||
콜드 비용을 시작 시 몰아, 이후 실 작업은 웜(빠름). 백그라운드로 돌려 잡 처리를 막지 않는다.
|
||
나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다.
|
||
시도당 타임아웃 필수 — 웜업은 search 중 어댑터 락을 쥐므로, 여기서 행하면 그 워커의
|
||
모든 실 검색이 락 대기로 함께 멈춘다(2026-07-10 부하테스트에서 15분 행 실측)."""
|
||
for ad in worker_adapters:
|
||
if ad.source not in _WARMUP_SOURCES:
|
||
continue
|
||
for attempt in range(tries):
|
||
try:
|
||
await asyncio.wait_for(ad.search("생수", limit=1), timeout=attempt_timeout)
|
||
LOG.i(f"[warmup:{ad.source}] 통과·쿠키 확보 (시도 {attempt + 1})")
|
||
if alerts is not None:
|
||
await alerts.check(f"warmup_{ad.source}", False, f"{ad.source} 크롤 정상")
|
||
break
|
||
except Exception as ex:
|
||
if attempt < tries - 1:
|
||
ad._rotate_ip(f"웜업 재시도({type(ex).__name__}) — 새 IP")
|
||
else:
|
||
msg = (f"[warmup:{ad.source}] {tries}회 모두 실패({type(ex).__name__}: {str(ex)[:80]}) — "
|
||
f"이 환경에서 {ad.source} 크롤이 막혔을 수 있습니다. "
|
||
f"호스트에서는 되는데 컨테이너에서만 막히는 사례가 있으니 실행 환경을 확인하세요")
|
||
LOG.w(msg)
|
||
if alerts is not None:
|
||
await alerts.check(f"warmup_{ad.source}", True, msg)
|
||
|
||
|
||
def _fatal_markers(adapters) -> list[str]:
|
||
"""'회전 무효' 마커 목록 — 알림이 이 마커만 세도록 모은다.
|
||
|
||
두 종류를 함께 센다: 어댑터가 HTML 에서 알아보는 구조적 차단 마커와, 연속 패턴으로
|
||
판정하는 환경 차단(서킷브레이커 트립, ENV_BLOCK_MARKER). 둘 다 IP 회전으로는 회복되지
|
||
않아 사람이 환경/설정을 고쳐야 한다는 점에서 같은 계열이다.
|
||
"""
|
||
out = [ENV_BLOCK_MARKER]
|
||
for ad in (adapters or []):
|
||
out += list(getattr(ad, "fatal_block_markers", ()) or ())
|
||
return sorted(set(out))
|
||
|
||
|
||
async def _port_pool_status(adapters) -> tuple[dict, tuple[int, int] | None]:
|
||
"""포트 장부 현황 → (게이트웨이별 상세, (최소 가용, 전체)).
|
||
|
||
게이트웨이가 둘(쿠팡=국가무지정 / 네이버=한국)이라 **가장 마른 쪽 기준(min)**으로 본다 —
|
||
한쪽만 고갈돼도 그 소스는 검색을 못 하므로 평균으로 덮으면 안 된다.
|
||
가용 수는 DB 장부에서 읽는다(프로세스 로컬 카운터는 다른 프로세스의 차단을 모른다).
|
||
"""
|
||
for ad in (adapters or []):
|
||
store = getattr(getattr(ad, "_proxy", None), "_store", None)
|
||
if store is None:
|
||
continue
|
||
try:
|
||
detail = await store.snapshot()
|
||
except Exception as ex:
|
||
LOG.d(f"[ops] 포트 장부 조회 실패(무시): {type(ex).__name__}")
|
||
return {}, None
|
||
if not detail:
|
||
return {}, None
|
||
worst = min(detail.values(), key=lambda d: d["available"])
|
||
return detail, (worst["available"], worst["total"])
|
||
return {}, None
|
||
|
||
|
||
async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0, adapters=None, alerts=None, ip_log=None):
|
||
"""워커 헬스 하트비트 + 임계 알림. 주기적으로 (1) 하트비트 파일 갱신(Docker HEALTHCHECK 가
|
||
행/좀비 워커 감지) (2) 큐/차단/DB풀/소스별 실패 지표 점검 → AlertManager 로 발화
|
||
(룰별 쿨다운으로 스팸 방지, 조건 해소 시 회복 알림)."""
|
||
hb_path = worker_config.heartbeat_file
|
||
th = alert_config # 임계값은 [AlertConfig] 섹션이 소스(docs/operations.md 표)
|
||
alerts = alerts or AlertManager(origin="worker")
|
||
ip_log = ip_log or IpSessionLog()
|
||
while not stop.is_set():
|
||
try:
|
||
with open(hb_path, "w") as f:
|
||
f.write(str(int(time.time()))) # 하트비트(mtime) — HEALTHCHECK 가 신선도 확인
|
||
except Exception:
|
||
pass
|
||
try:
|
||
snap = await queue.ops()
|
||
snap["blocks_1h"] = await bot_log.recent_count(60)
|
||
pool = DB_SESSION_MNG.pool_status()
|
||
snap["pool_pct"] = pool["pct"]
|
||
await alerts.check("dead", snap["dead_1h"] >= th.dead_1h, f"DEAD 1h={snap['dead_1h']}", snap)
|
||
await alerts.check("blocks", snap["blocks_1h"] >= th.blocks_1h, f"차단 1h={snap['blocks_1h']}", snap)
|
||
# 구조적 차단(회전 무효) — 1건만 나와도 알린다. 방치하면 그 소스는 계속 0건이다.
|
||
fatal = _fatal_markers(adapters)
|
||
if fatal:
|
||
snap["fatal_blocks_1h"] = await bot_log.recent_count_by_marker(fatal, 60)
|
||
await alerts.check("fatal_block", snap["fatal_blocks_1h"] > 0,
|
||
f"구조적 차단 1h={snap['fatal_blocks_1h']} (마커 {fatal}) — "
|
||
f"IP 회전으로 회복 불가. 게이트웨이 국가 설정([DecodoConfig].kr_host) 확인", snap)
|
||
await alerts.check("queue_lag", snap["oldest_pending_sec"] >= th.queue_lag_sec, f"큐지연={snap['oldest_pending_sec']}s", snap)
|
||
await alerts.check("stuck", snap["stuck_running"] > 0, f"stuck={snap['stuck_running']}", snap)
|
||
await alerts.check("db_pool", pool["pct"] >= th.pool_pct,
|
||
f"DB 풀 포화 {pool['pct']}% (checked_out {pool['checked_out']}/{pool['capacity']})", snap)
|
||
await alerts.check("deadline", snap["deadline_1h"] >= th.deadline_1h,
|
||
f"잡 데드라인 강제종료 1h={snap['deadline_1h']} — 크롤 행 반복 신호", snap)
|
||
await alerts.check("cost", snap["cost_1h_usd"] >= th.cost_1h_usd,
|
||
f"검색원가 1h=${snap['cost_1h_usd']} — 비용 폭주(리소스차단 풀림·재시도 루프) 점검", snap)
|
||
# 가용 프록시 포트 고갈 — 쿨다운 격리 누적. blocks_1h 보다 먼저 우는 대규모 차단 조기 신호.
|
||
pool, ports = await _port_pool_status(adapters)
|
||
if ports:
|
||
avail, total = ports
|
||
snap["proxy_ports_avail"], snap["proxy_ports_total"] = avail, total
|
||
if pool:
|
||
snap["port_pool"] = pool # {게이트웨이: {held, resting, cooling, available, total}}
|
||
await alerts.check("proxy_ports_low", avail * 100 <= total * th.ports_low_pct,
|
||
f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호"
|
||
+ (f" · {pool}" if pool else ""), snap)
|
||
# 예산 누수 — 요청 예산을 지켰는데도 차단된 IP 세션 발생.
|
||
# 처방은 '몇 번째 요청에서 막혔나'로 갈린다:
|
||
# ip_req#1 위주 → 새 IP 첫 요청부터 차단 = IP 평판 문제. 예산을 낮춰도 소용없다.
|
||
# ip_req#2~ 위주 → 같은 IP 로 너무 많이 긁은 것 = 예산 하향이 유효.
|
||
# 예산 하향만 권하면 오진을 부른다(2026-07-28 배포서버 조사에서 전량 ip_req#1 이었다).
|
||
block_sessions = (await ip_log.recent_stats(360)).get("block", 0)
|
||
snap["block_sessions_6h"] = block_sessions
|
||
await alerts.check(
|
||
"budget_leak", block_sessions >= th.block_sessions_6h,
|
||
f"예산 회전에도 차단된 IP 세션 6h={block_sessions} — "
|
||
f"bot_detection.ip_request_no 분포 확인(1 위주면 IP 평판/프록시 대역, 2 이상이면 ip_request_budget 하향)",
|
||
snap,
|
||
)
|
||
# 소스별 장기 실패 — 최근 30분간 시도는 있는데 성공이 0건(쿼터 소진·셀렉터 드리프트·전면 차단 신호)
|
||
per_source: dict[str, list[int]] = {}
|
||
for ad in (adapters or []):
|
||
tries, ok = ad.recent_stats(1800)
|
||
agg = per_source.setdefault(ad.source, [0, 0])
|
||
agg[0] += tries
|
||
agg[1] += ok
|
||
for src, (tries, ok) in per_source.items():
|
||
await alerts.check(f"source_fail:{src}", tries >= th.source_fail_30m and ok == 0,
|
||
f"{src} 최근 30분 {tries}회 시도·성공 0건", snap)
|
||
except Exception as ex:
|
||
LOG.e_no_callstack(f"[ops-monitor] {type(ex).__name__}: {ex}")
|
||
try:
|
||
await asyncio.wait_for(stop.wait(), timeout=interval)
|
||
except asyncio.TimeoutError:
|
||
pass
|
||
|
||
|
||
async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: float = 30.0,
|
||
close_timeout: float = 60.0):
|
||
"""유휴 브라우저 정리 루프 — 일정 시간 검색 없는 어댑터의 Chrome 을 닫아 메모리를 회수한다.
|
||
쿠키는 user_data_dir 에 남아, 다음 검색 때 재기동해도 (같은 IP면) 웜 유지.
|
||
순차 순회라 close 1건에도 타임아웃을 건다 — 한 어댑터의 close 행이 루프 전체를 멈춰
|
||
다른 워커의 브라우저까지 못 닫게 되는 것을 실측(2026-07-10 부하테스트)했다."""
|
||
while not stop.is_set():
|
||
try:
|
||
await asyncio.wait_for(stop.wait(), timeout=interval)
|
||
except asyncio.TimeoutError:
|
||
pass
|
||
for ad in adapters:
|
||
close_if_idle = getattr(ad, "close_if_idle", None)
|
||
if close_if_idle is None: # 네이버(httpx) 등 브라우저 없는 어댑터는 정리 대상 아님
|
||
continue
|
||
try:
|
||
await asyncio.wait_for(close_if_idle(idle_sec), timeout=close_timeout)
|
||
except asyncio.TimeoutError:
|
||
LOG.w(f"[browser-reaper] {getattr(ad, 'source', '?')} 정리 {close_timeout:.0f}s 초과 — 취소·스킵(다음 틱 재시도)")
|
||
except Exception as ex:
|
||
LOG.e_no_callstack(f"[browser-reaper] 정리 실패(무시): {ex}")
|
||
|
||
|
||
async def main(concurrency: int = 1):
|
||
queue = JobQueue()
|
||
neg_cache, history = NegativeCache(), PriceHistory() # DB 기반 — 워커 공유 안전
|
||
has_openai = bool(openai_config.api_key)
|
||
|
||
# 시작 프리플라이트: DECODO 게이트가 살아있는지(인증) 대표 프록시로 1회 확인. 포트는 워커별로 각자 잡음.
|
||
probe = DecodoProxy()
|
||
LOG.i(f"DECODO 프록시: {'ON(sticky ' + str(probe.session_minutes) + '분 회전)' if probe.enabled else 'OFF(미설정)'}")
|
||
if probe.enabled:
|
||
egress_ip, egress_port = await probe.healthcheck()
|
||
LOG.i(f"DECODO 프리플라이트 OK — egress IP {egress_ip} (port {egress_port})") if egress_ip \
|
||
else LOG.w("DECODO 프리플라이트 실패 — 살아있는 포트를 못 찾음(런타임 회전으로 재시도)")
|
||
fb = _enabled_fallbacks()
|
||
LOG.i(f"AI(판정+검색어생성): {'ON' if has_openai else 'OFF(키 없음)'} · "
|
||
f"오픈마켓 폴백: {', '.join(fb) if fb else 'OFF(기본 — LPS_FALLBACKS 로 활성화)'}")
|
||
|
||
stop = asyncio.Event()
|
||
listeners: list[JobListener] = []
|
||
tasks: list[asyncio.Task] = []
|
||
bg_tasks: list[asyncio.Task] = [] # 웜업 등 백그라운드(짧게 끝남, gather 대상 아님)
|
||
all_adapters = []
|
||
|
||
# ── graceful shutdown: SIGINT(Ctrl+C)/SIGTERM(docker stop) → stop 이벤트 ──
|
||
# asyncio.run 기본 동작(SIGINT=메인 태스크 즉시 cancel)은 하던 잡을 도중에 끊어
|
||
# RUNNING 인 채 lease 만료(120s)까지 묶어둔다. 대신 stop 을 set 해 "새 잡은 안 받고,
|
||
# 하던 잡은 마무리"로 종료한다. 같은 신호를 한 번 더 받으면 강제 종료(태스크 취소).
|
||
def _request_stop(sig_name: str):
|
||
if not stop.is_set():
|
||
LOG.i(f"{sig_name} 수신 — graceful 종료: 새 잡 중단, 하던 잡 마무리 (한 번 더 = 강제 종료)")
|
||
stop.set()
|
||
for t in bg_tasks: # 웜업은 선택 작업 — 즉시 취소해 어댑터 락을 비운다
|
||
t.cancel()
|
||
else:
|
||
LOG.w(f"{sig_name} 재수신 — 강제 종료(실행 중 잡은 lease 만료 후 reaper 가 재큐)")
|
||
for t in tasks:
|
||
t.cancel()
|
||
|
||
loop = asyncio.get_running_loop()
|
||
for sig in (signal.SIGINT, signal.SIGTERM):
|
||
loop.add_signal_handler(sig, _request_stop, sig.name)
|
||
|
||
# 잡 1건 데드라인 — 정상 검색은 폴백 포함 수분 내 끝난다(실측 15~22s). 크롤 행 실측(15분) 대비 상한.
|
||
job_deadline = worker_config.job_deadline_sec
|
||
# 알림은 웜업과 ops 모니터가 공유한다 — 같은 룰이 양쪽에서 중복 발화하지 않도록.
|
||
alerts = AlertManager(origin="worker")
|
||
# 포트(=IP 세션) 장부는 DB 다 — 프로세스가 늘어도 모두 같은 장부를 본다.
|
||
# 기동 시 게이트웨이별 포트 행을 보장한다(이미 있으면 상태 유지, 덮지 않음).
|
||
store = PortLeaseStore()
|
||
for gw in {decodo_config.host, decodo_config.kr_host or decodo_config.host}:
|
||
if gw:
|
||
await store.ensure_ports(gw, decodo_config.port_start, decodo_config.port_end)
|
||
for i in range(concurrency):
|
||
handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history, store=store)
|
||
all_adapters += worker_adapters
|
||
bg_tasks.append(asyncio.create_task(_warmup_worker(worker_adapters, alerts=alerts))) # 쿠키 선점 + 크롤 프리플라이트
|
||
listener = JobListener()
|
||
await listener.start()
|
||
listeners.append(listener)
|
||
worker = Worker(f"worker-{i}", queue, handler, job_deadline_sec=job_deadline)
|
||
tasks.append(asyncio.create_task(worker.run(listener, stop)))
|
||
|
||
tasks.append(asyncio.create_task(run_reaper(queue, stop)))
|
||
tasks.append(asyncio.create_task(run_browser_reaper(all_adapters, stop))) # 유휴 브라우저 정리
|
||
tasks.append(asyncio.create_task(run_ops_monitor(queue, BotDetectionLog(), stop, adapters=all_adapters, alerts=alerts))) # 하트비트 + 임계 알림
|
||
LOG.i(f"LPS 워커 {concurrency}개 + reaper + 브라우저정리 + ops모니터(하트비트/알림) 기동 (워커별 세트 · 상품 {concurrency}개 동시)")
|
||
|
||
# 종료 유예: stop 후 하던 잡이 이 시간 안에 끝나면 자연 종료, 초과하면 강제 취소.
|
||
# docker stop 을 쓰면 compose 의 stop_grace_period 를 이보다 길게 잡아야 SIGKILL 전에 마무리된다.
|
||
grace = worker_config.shutdown_grace_sec
|
||
gathered = asyncio.gather(*tasks)
|
||
stop_waiter = asyncio.create_task(stop.wait())
|
||
try:
|
||
await asyncio.wait({gathered, stop_waiter}, return_when=asyncio.FIRST_COMPLETED)
|
||
if gathered.done():
|
||
gathered.result() # 워커/리퍼가 예외로 죽은 경우 → 전파(finally 가 정리 후 종료)
|
||
else:
|
||
# 종료 신호 경로 — 워커 루프들이 stop 을 보고 하던 잡을 마친 뒤 스스로 끝나길 기다린다
|
||
try:
|
||
await asyncio.wait_for(gathered, timeout=grace)
|
||
LOG.i("graceful 종료 — 모든 워커가 하던 잡을 마무리함")
|
||
except asyncio.TimeoutError:
|
||
LOG.w(f"종료 유예 {grace:.0f}s 초과 — 남은 태스크 강제 취소(잡은 lease 만료 후 재큐)")
|
||
except asyncio.CancelledError: # 신호 재수신(강제 종료)로 태스크가 취소된 경우
|
||
LOG.w("강제 종료 — 남은 리소스 정리 후 종료")
|
||
finally:
|
||
stop.set()
|
||
stop_waiter.cancel()
|
||
for t in (*tasks, *bg_tasks):
|
||
t.cancel()
|
||
# 취소 완주를 기다린 뒤 정리 — 실행 중 태스크가 브라우저/커넥션을 쓰는 채로 닫지 않게
|
||
await asyncio.gather(gathered, stop_waiter, *bg_tasks, return_exceptions=True)
|
||
for listener in listeners:
|
||
try:
|
||
await listener.close()
|
||
except Exception as ex:
|
||
LOG.e_no_callstack(f"[shutdown] 리스너 정리 실패(무시): {ex}")
|
||
for adapter in all_adapters: # 항목별 격리 — 하나가 실패해도 나머지 Chrome 은 닫는다
|
||
try:
|
||
await adapter.close()
|
||
except Exception as ex:
|
||
LOG.e_no_callstack(f"[shutdown] {getattr(adapter, 'source', '?')} 정리 실패(무시): {ex}")
|
||
LOG.i("LPS 워커 종료 완료")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
# 동시성은 [WorkerConfig].concurrency 가 소스 — WORKER_CONCURRENCY env 는 실행 스크립트의
|
||
# 대화형 입력 전용 임시 override(설정 관리는 toml 하나로, 2026-07-13 협의).
|
||
_conc = int(os.environ.get("WORKER_CONCURRENCY", "0")) or worker_config.concurrency
|
||
asyncio.run(main(_conc))
|