o2o-negosium-original/lps/worker_main.py
민헌 c955a8f367 feat(lps): 프록시 전송오류 IP회전 + 시작 프리플라이트 + DECODO/컴포넌트별 비용
프록시(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>
2026-07-09 15:33:37 +09:00

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"))))