프록시(DECODO) 포트/IP 사망(407/ERR_TUNNEL/ERR_HTTP_RESPONSE_CODE_FAILURE)이
봇차단과 구분 없이 예외로 튕겨 같은 죽은 포트로 재시도만 하다 DEAD 되던 문제를 고친다.
- browser_base: is_proxy_error(순수함수) + search 루프에서 프록시 전송오류 시 IP 회전 재시도
(max_proxy_retries=2). 봇감지 회전과 통합. uses_proxy 프로퍼티.
- 쿠팡 어댑터를 BrowserSearchAdapter 로 통합 — 중복 machinery 제거, 회전 로직 한 곳에서 공유
(detect_block 순수함수는 유지, 테스트 호환).
- proxy.healthcheck(): 시작 프리플라이트 — 살아있는 포트 선점 + egress IP 로그(빠른 실패·가시성).
worker_main 기동 시 호출.
- 비용: DecodoConfig.cost_per_gb 추가. metrics 에 proxy_bytes(네이버 직접 제외) + 컴포넌트별
cost{ai_usd, proxy_usd, total_usd}. FE 원가 타일에 AI/DECODO 분해·프록시 바이트.
- 테스트: is_proxy_error 8종 + 비용 분해 1종.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
98 lines
4.6 KiB
Python
98 lines
4.6 KiB
Python
# 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"))))
|