o2o-negosium-original/lps/worker_main.py
민헌 23c1c48634 feat(lps): 네이버 최저가를 모바일 크롤로 복구(오픈API 종료 대체)
shop.json 이 2026-07-31 종료(404 SE05)되고 NCP API HUB 에도 승계되지 않아
가격을 얻을 공식 경로가 사라졌다 → 쿠팡과 같은 스택(patchright+실제 Chrome)으로 크롤 전환.

경로: msearch.shopping.naver.com (PC 는 405/418 로 막힘). 7/9 스파이크 때 모바일은
로그인 리다이렉트였는데 그 사이 열렸다.

통과 조건 3개 — 하나라도 빠지면 WTM 캡차(실측):
- **한국 IP**: 해외 residential 은 즉시 하드차단(2.6KB) → kr.decodo.com 게이트웨이
  ([DecodoConfig].kr_host, DecodoProxy(host=...) 로 주입. 쿠팡은 기존 월드와이드 유지)
- **ko-KR 로케일/시간대**: KR IP + en-US 조합을 봇으로 본다
  (BrowserSearchAdapter.context_options 훅 추가)
- **리소스 차단 금지**: route 를 걸면 즉시 캡차. image/media/font 만 막아도 동일 →
  '무엇을 막느냐'가 아니라 요청 가로채기 자체가 탐지 신호. 대신 검색당 ~3MB(~$0.009)

파서는 '정확한 상품의 최저가'를 기준으로 취사선택한다:
- 광고/슈퍼적립/브랜드블록 카드 제외(멤버십·쿠폰 조건부 가격)
- 쿠폰할인가를 price 로 쓰지 않음(조건부라 실구매가보다 싸게 잡힘)
- 가격비교('최저 N원') 카드는 유지하고 mall_name="네이버"(옛 lprice 와 같은 의미)
- **배송비 확보** — 옛 오픈API 는 필드 자체가 없어 전 소스 None 이었다
- 가격 함정 3종 회귀 테스트: 단위가격(548원)·가격노드 안의 배송비(3,900원)·정상가/할인율

source 는 "naver" 유지 — price_history.naver_lowest·MALL_BY_SOURCE·프론트 그래프 계약이
구현(API→크롤) 교체와 무관하게 살아야 한다.

테스트 12건 추가(축약 픽스처 + 합성 함정) · 전체 182 passed.
2026-08-05 09:05:38 +09:00

331 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 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.price_history import PriceHistory
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):
"""워커 1개의 자립 세트(브라우저 어댑터·AI·핸들러)를 만든다.
프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP)로 진짜 병렬을 보장한다."""
# 워커별 프록시(다른 포트=다른 IP). 100포트를 워커 수로 균등 분할해 시작점을 벌린다.
proxy = DecodoProxy()
# 네이버 전용 프록시 — 같은 계정·같은 포트지만 **한국 타깃 게이트웨이**를 쓴다.
# 해외 residential IP 로는 msearch 가 즉시 하드차단된다(2026-08-04 실측). kr_host 가 비면
# 기본 게이트웨이로 떨어지고, 그때는 네이버가 막힐 수 있다(로그의 BOT-DETECTED 로 드러남).
kr_proxy = DecodoProxy(host=decodo_config.kr_host or None)
if proxy.enabled and concurrency > 1:
n = proxy.port_end - proxy.port_start + 1
proxy.seed_offset(i * max(1, n // concurrency))
kr_proxy.seed_offset(i * max(1, n // concurrency))
bot_log = BotDetectionLog()
ip_log = IpSessionLog()
suffix = f"_w{i}" if concurrency > 1 else ""
def _pf(source): # 워커별 Chrome 프로필 경로(중복 실행 시 ProcessSingleton 충돌 방지)
# [WorkerConfig].profile_dir 를 영속 볼륨으로 두면 재시작해도 cf_clearance 등 쿠키 유지(재웜업 회피).
return f"{worker_config.profile_dir}/lps_{source}{suffix}"
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,
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())
async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float = 60.0):
"""워커의 챌린지 소스(Turnstile/Akamai)를 미리 풀어 쿠키(cf_clearance 등)를 확보한다.
콜드 비용을 시작 시 몰아, 이후 실 작업은 웜(빠름). 백그라운드로 돌려 잡 처리를 막지 않는다.
나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다.
시도당 타임아웃 필수 — 웜업은 search 중 어댑터 락을 쥐므로, 여기서 행하면 그 워커의
모든 실 검색이 락 대기로 함께 멈춘다(2026-07-10 부하테스트에서 15분 행 실측)."""
for ad in worker_adapters:
if ad.source not in ("gmarket", "auction", "coupang"):
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})")
break
except Exception as ex:
if attempt < tries - 1:
ad._rotate_ip(f"웜업 재시도({type(ex).__name__}) — 새 IP")
else:
LOG.w(f"[warmup:{ad.source}] {tries}회 실패(첫 잡에서 재시도): {type(ex).__name__}")
def _proxy_ports_snapshot(adapters) -> tuple[int, int] | None:
"""워커 프록시들의 가용 포트 현황 — (최소 가용 수, 전체 포트 수). 프록시 미사용이면 None.
쿨다운 맵은 프록시 인스턴스(워커)별이라 가장 소진된 워커 기준(min)으로 본다."""
proxies = {}
for ad in (adapters or []):
p = getattr(ad, "_proxy", None)
if p is not None and p.enabled:
proxies[id(p)] = p
if not proxies:
return None
any_p = next(iter(proxies.values()))
total = any_p.port_end - any_p.port_start + 1
return min(p.available_ports() for p in proxies.values()), total
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)
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 보다 먼저 우는 대규모 차단 조기 신호.
ports = _proxy_ports_snapshot(adapters)
if ports:
avail, total = ports
snap["proxy_ports_avail"], snap["proxy_ports_total"] = avail, total
await alerts.check("proxy_ports_low", avail * 100 <= total * th.ports_low_pct,
f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호", 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
for i in range(concurrency):
handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history)
all_adapters += worker_adapters
bg_tasks.append(asyncio.create_task(_warmup_worker(worker_adapters))) # 챌린지 쿠키 선점(백그라운드)
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))) # 하트비트 + 임계 알림
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))