# LPS 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일). # python worker_main.py # WORKER_CONCURRENCY=2 python worker_main.py # # 쿠팡 검색은 브라우저(Chrome)라 무겁고 컨텍스트당 직렬이므로 기본 동시성은 1. # 여러 브라우저로 늘리려면 CoupangAdapter 인스턴스를 워커마다 따로 준다. import asyncio import os from common.logger import LOG from config.server_configs import web_server_config, openai_config, decodo_config from crud.job_crud import JobQueue from crud.negative_cache import NegativeCache from crud.bot_detection import BotDetectionLog from crud.price_history import PriceHistory from services.search.proxy import DecodoProxy from services.search.coupang.adapter import CoupangAdapter from services.search.naver.adapter import NaverAdapter 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") async def main(concurrency: int = 1): queue = JobQueue() # 쿠팡(브라우저, 무거움) + 네이버(오픈API, 가벼움) 동시 검색 → 병합 최저가 # DECODO 프록시: .env 에 값 있으면 쿠팡만 sticky+주기적 회전으로 경유(없으면 직접 연결) proxy = DecodoProxy() LOG.i(f"DECODO 프록시: {'ON(sticky ' + str(proxy.session_minutes) + '분 회전)' if proxy.enabled else 'OFF(미설정)'}") # 시작 프리플라이트: 살아있는 프록시 포트를 선점하고 egress IP 를 로그로 남긴다(빠른 실패·가시성). # residential IP 는 실행 중에도 죽으므로 실제 회복은 런타임 IP 회전(전송오류·봇감지)이 담당. if proxy.enabled: egress_ip, egress_port = await proxy.healthcheck() if egress_ip: LOG.i(f"DECODO 프리플라이트 OK — egress IP {egress_ip} (port {egress_port})") else: LOG.w("DECODO 프리플라이트 실패 — 살아있는 포트를 못 찾음(런타임 회전으로 재시도)") # 봇 감지 시: 감지 기록(DB) + 새 IP 로 회전 후 재시도 bot_log = BotDetectionLog() adapters = { "coupang": CoupangAdapter(headless=False, proxy=proxy, on_detect=bot_log.record), "naver": NaverAdapter(), } # 오픈마켓 폴백 크롤러: 네이버가 그 몰을 커버 못 했을 때만 lazy 하게 실사이트 크롤(브라우저는 첫 사용 시 기동). # G마켓·옥션(ESM '잠시만' 챌린지) + 11번가(PC). 리소스차단 OFF(렌더/챌린지 보호)는 어댑터 기본값. fallback_adapters = { "gmarket": EsmAdapter("gmarket", headless=False, proxy=proxy, on_detect=bot_log.record), "auction": EsmAdapter("auction", headless=False, proxy=proxy, on_detect=bot_log.record), "st11": ElevenStAdapter(headless=False, proxy=proxy, on_detect=bot_log.record), } LOG.i(f"오픈마켓 폴백 크롤: {', '.join(fallback_adapters)} (네이버 미커버 몰만)") # OpenAI 키 있으면 '같은 상품' AI 판정 + 재검색어 생성 활성화 has_openai = bool(openai_config.api_key) judge = SimilarityJudge() if has_openai else None keyword_gen = KeywordGenerator() if has_openai else None LOG.i(f"AI(판정+검색어생성): {'ON' if has_openai else 'OFF(키 없음)'}") handler = build_search_handler( adapters, judge=judge, keyword_gen=keyword_gen, neg_cache=NegativeCache(), history=PriceHistory(), fallback_adapters=fallback_adapters, ai_model=openai_config.model, proxy_cost_per_gb=decodo_config.cost_per_gb, ) stop = asyncio.Event() listeners: list[JobListener] = [] tasks: list[asyncio.Task] = [] for i in range(concurrency): listener = JobListener() await listener.start() listeners.append(listener) worker = Worker(f"worker-{i}", queue, handler) tasks.append(asyncio.create_task(worker.run(listener, stop))) tasks.append(asyncio.create_task(run_reaper(queue, stop))) LOG.i(f"LPS 워커 {concurrency}개 + reaper 기동") try: await asyncio.gather(*tasks) finally: stop.set() for listener in listeners: await listener.close() for adapter in list(adapters.values()) + list(fallback_adapters.values()): await adapter.close() if __name__ == "__main__": asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))