o2o-negosium-original/lps/worker_main.py
민헌 495e16b951 feat(lps): IP 를 태우지 말고 한계 직전까지 쓰고 쉬게 — 네이버 예산 상향 + 휴식 개념
기존엔 네이버가 쿠팡 기준 예산(3회)을 그대로 썼다. 실측하니 체급이 다르다:
같은 KR IP 로 **12회 연속 검색까지 무차단**(IP 4개 전부 한계 미도달). 3회로 돌리면
불필요하게 4배 자주 회전해 KR 풀만 빨리 소모하고 회전마다 브라우저 재기동(~20s)이 붙는다.
→ [DecodoConfig].naver_ip_request_budget = 10 (실측 12 에 여유). 쿠팡은 3 유지.

그리고 선제 회전에 빠져 있던 조각을 채웠다 — **휴식(rest)**:
예산 도달로 놓은 포트를 곧바로 다른 워커가 집으면 그 IP 의 요청률이 도로 올라가
예산의 의미가 사라진다. release(rest_sec=...) 로 sticky 수명만큼 쉬게 한다.
차단으로 태우는 burn(30분)과는 별개 상태다:
  휴식  탄 게 아님 · 짧음 · 소진 시 가장 먼저 회수
  쿨다운 차단당함 · 김 · 휴식보다 나중에 회수
회전 종류(kind)를 browser_base → DecodoProxy.rotate(kind) 로 전달해 budget 일 때만 휴식을 건다.

라이브 검증(예산 3으로 낮춰 관찰): 6회 검색 = IP 2개만 사용, 3회마다 선제 회전,
놓은 포트는 휴식 1 · 쿨다운 0 · 차단 0. 즉 IP 를 태우지 않고 로테이션만으로 돌아간다.

테스트 6건 추가(휴식 재사용 금지·만료 복귀·burn 우선·회수 우선순위·budget vs block),
전체 202 passed. _MockProxy.rotate 가 kind 를 받도록 갱신.
2026-08-05 11:16:37 +09:00

354 lines
22 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.port_registry import PortRegistry
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, registry=None):
"""워커 1개의 자립 세트(브라우저 어댑터·AI·핸들러)를 만든다.
프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP)로 진짜 병렬을 보장한다."""
# 워커별 프록시(다른 포트=다른 IP). 100포트를 워커 수로 균등 분할해 시작점을 벌린다.
# 두 프록시는 **게이트웨이가 다르다** — 쿠팡=국가 무지정, 네이버=한국 전용.
# 해외 residential IP 로는 네이버 msearch 가 즉시 하드차단된다(2026-08-05 실측:
# 같은 포트 10091 에서 gate=차단 2,641B / kr=정상 14건). kr_host 가 비면 gate 로
# 떨어지고 그때는 네이버가 막힌다(로그의 [naver][BOT-DETECTED] 로 드러남).
# 포트는 공용 registry 가 배타 임대해 준다 — 워커끼리, 그리고 두 소스끼리 같은 IP 를
# 동시에 쓰거나 태운 IP 를 곧바로 재사용하는 걸 막는다.
proxy = DecodoProxy(registry=registry, owner=f"coupang-w{i}")
kr_proxy = DecodoProxy(host=decodo_config.kr_host or None, registry=registry, owner=f"naver-w{i}")
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,
ip_request_budget=decodo_config.naver_ip_request_budget,
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
def _port_pool_detail(adapters) -> dict:
"""게이트웨이별 상세(보유·쿨다운·가용·소유자) — 알림 페이로드에 실어 원인 파악을 돕는다.
registry 를 쓰지 않는 구성(단독 프록시)이면 빈 dict."""
for ad in (adapters or []):
reg = getattr(getattr(ad, "_proxy", None), "_registry", None)
if reg is not None:
return reg.snapshot()
return {}
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
pool = _port_pool_detail(adapters)
if pool:
snap["port_pool"] = pool # {게이트웨이: {held, cooling, available, owners}}
await alerts.check("proxy_ports_low", avail * 100 <= total * th.ports_low_pct,
f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호"
+ (f" · {pool}" if pool else ""), 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
# 프로세스 전체가 공유하는 포트(=IP 세션) 중재자 — 워커·소스가 늘어도 여기 하나만 본다.
registry = PortRegistry(decodo_config.port_start, decodo_config.port_end)
for i in range(concurrency):
handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history, registry=registry)
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))