o2o-negosium-original/lps/worker_main.py
민헌 69f6641c27 feat(lps): 검색 1건 원가 계측 — AI 토큰·비용 + 시간 + 크롤 트래픽
검색이 소모하는 리소스/비용/시간을 잡 단위로 집계해 result.metrics 로 적재(API/FE 노출).
지금까진 타임스탬프만 있고 실제 비용 동인(AI 토큰·대역폭)은 버려지고 있었다.

- services/metrics.SearchMetrics: duration_ms + ai(calls/tokens/est_cost_usd, gpt-4o-mini 단가)
  + crawl(fetches/html_bytes/malls_crawled) + source_ms
- AI 클라이언트: resp.usage 를 last_usage 로 노출(그동안 폐기하던 토큰)
- 어댑터: last_bytes(처리 HTML 바이트) 노출 — naver/coupang/browser_base 공통
- handler: 각 fetch 타이밍+바이트, AI 호출 토큰을 metrics 로 누적 → 결과에 스냅샷
- FE: 작업 카드에 원가 4타일(소요/AI비용/토큰/크롤 트래픽)
- 테스트 2종. ⚠️ html_bytes 는 대역폭 근사(오픈마켓 리소스 미차단분 제외=하한), CDP 정확화는 백로그

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 15:14:35 +09:00

89 lines
4.0 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
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(미설정)'}")
# 봇 감지 시: 감지 기록(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,
)
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"))))