diff --git a/lps/services/search/port_registry.py b/lps/services/search/port_registry.py new file mode 100644 index 0000000..cf0cd6c --- /dev/null +++ b/lps/services/search/port_registry.py @@ -0,0 +1,124 @@ +"""DECODO 포트(=IP 세션) 중재자 — 누가 어떤 IP 를 쓰고 있고, 어떤 IP 가 탔는지 한곳에서 관리한다. + +## 왜 필요한가 + +DECODO 는 포트 1개 = sticky 세션 1개다. 그런데 우리는 **한 계정으로 두 종류의 풀**을 쓴다: + + 쿠팡 gate.decodo.com:10001-10100 국가 무지정(실측: VN·MY·BD·ID·KZ·IN…) + 네이버 kr.decodo.com:10001-10100 한국 전용(실측: LG U+·KT·SK브로드밴드) + +네이버는 해외 IP 를 즉시 하드차단하고(2.6KB), 쿠팡은 한국 IP 가 필수는 아니다. 그래서 +소스마다 게이트웨이가 갈린다. **같은 포트 번호라도 게이트웨이가 다르면 IP 가 다르다** +(실측: port 10061 → gate=103.99.27.55(ID) / kr=121.180.128.2(KR)) — 그래서 키는 (host, port) 다. + +여기가 없으면 생기는 문제: + 1. 워커 N개가 같은 게이트웨이를 쓰면 시간이 지나며 같은 포트로 수렴해 **한 IP 에 요청이 몰린다** + (seed_offset 은 시작점만 벌릴 뿐, 회전하다 보면 겹친다) → 그 IP 가 빨리 탄다 + 2. 차단으로 태운 포트를 **다른 소유자가 곧바로 다시 집는다** — 쿨다운이 프록시 인스턴스별이라서 + 3. kr_host 를 비워 네이버가 gate 로 폴백하면 쿠팡과 같은 풀을 나눠 쓰게 된다. 이때는 + 포트 번호가 곧 같은 IP 라 1·2 번이 실제 IP 충돌로 이어진다 + +## 설계 + +- 소유권은 **배타적**이다. 한 (host, port) 는 동시에 한 소유자만 쥔다. +- 쿨다운(burn)은 **전역**이다. 누가 태웠든 만료 전까지 아무도 못 쓴다. +- 리스는 session_minutes 로 만료된다(sticky 세션 수명과 맞춤). 만료되면 자동 반납된다. +- 전 포트가 막히면 **가장 빨리 풀리는 포트**를 내준다(가용성 우선 — 멈추는 것보다 낫다). + +프로세스 안에서만 공유한다. 워커 동시성은 한 프로세스 안의 asyncio 태스크라 이걸로 충분하고, +프로세스를 늘리면 그때 DB(ip_session)로 올려야 한다 — 지금 그 복잡도를 미리 지불하지 않는다. +""" + +import time + +from common.logger import LOG + + +class PortRegistry: + def __init__(self, port_start: int, port_end: int): + self.port_start = port_start + self.port_end = port_end + self._held: dict[tuple[str, int], tuple[str, float]] = {} # (host,port) → (owner, 취득시각) + self._burned: dict[tuple[str, int], float] = {} # (host,port) → 쿨다운 만료(monotonic) + + @property + def size(self) -> int: + return self.port_end - self.port_start + 1 + + # ── 할당 ──────────────────────────────────────────────────────────── + def acquire(self, host: str, owner: str, start_at: int = 0, lease_sec: float = 600) -> int: + """(host, *) 중 비어있고 쿨다운 아닌 포트를 배타 배정. start_at 은 탐색 시작 오프셋(워커 분산용). + + 전부 막혀 있으면 가장 먼저 풀리는 포트를 뺏어서라도 내준다 — 검색이 멈추는 것보다 낫다. + """ + self._prune(host) + n = self.size + for k in range(n): + port = self.port_start + ((start_at + k) % n) + key = (host, port) + if key in self._burned or key in self._held: + continue + self._held[key] = (owner, time.monotonic()) + return port + + # 전 포트 소진 — 쿨다운 만료가 가장 임박한 포트를 회수해 넘긴다. + victim = min(self._burned, key=self._burned.get, default=None) + if victim is None: # 쿨다운은 없는데 전부 점유 중 = 소유자가 포트 수보다 많다 + victim = min(self._held, key=lambda k: self._held[k][1]) + LOG.w(f"[port] {host} 가용 포트 없음 — {victim[1]} 회수해 {owner} 에 배정 " + f"(보유 {len(self._held)} · 쿨다운 {len(self._burned)} / {n})") + self._burned.pop(victim, None) + self._held[victim] = (owner, time.monotonic()) + return victim[1] + + def release(self, host: str, port: int, owner: str): + """리스 반납. 소유자가 다르면 무시한다(뒤늦은 반납이 남의 리스를 깨지 않도록).""" + key = (host, port) + cur = self._held.get(key) + if cur and cur[0] == owner: + del self._held[key] + + def burn(self, host: str, port: int, cooldown_sec: float, owner: str = "", reason: str = ""): + """차단·소진된 포트를 전역 격리. 보유 중이면 함께 반납한다.""" + key = (host, port) + self._burned[key] = time.monotonic() + cooldown_sec + self._held.pop(key, None) + LOG.i(f"[port] {host}:{port} 쿨다운 {int(cooldown_sec)}s ({owner}{' · ' + reason if reason else ''}) " + f"— 가용 {self.available(host)}/{self.size}") + + def lease_expired(self, host: str, port: int, owner: str, lease_sec: float) -> bool: + """리스가 sticky 수명을 넘겼는지(넘겼으면 호출부가 반납 후 재취득).""" + cur = self._held.get((host, port)) + if cur is None or cur[0] != owner: + return True + return (time.monotonic() - cur[1]) > lease_sec + + # ── 관측 ──────────────────────────────────────────────────────────── + def available(self, host: str) -> int: + self._prune(host) + used = sum(1 for (h, _) in self._held if h == host) + cooling = sum(1 for (h, _) in self._burned if h == host) + return self.size - used - cooling + + def snapshot(self) -> dict: + """게이트웨이별 현황 — ops 모니터/알림용. {host: {held, cooling, available, owners}}""" + self._prune() + hosts = {h for (h, _) in self._held} | {h for (h, _) in self._burned} + out = {} + for h in sorted(hosts): + owners: dict[str, int] = {} + for (hh, _), (owner, _) in self._held.items(): + if hh == h: + owners[owner] = owners.get(owner, 0) + 1 + out[h] = { + "held": sum(1 for (hh, _) in self._held if hh == h), + "cooling": sum(1 for (hh, _) in self._burned if hh == h), + "available": self.available(h), + "total": self.size, + "owners": owners, + } + return out + + def _prune(self, host: str | None = None): + now = time.monotonic() + self._burned = {k: t for k, t in self._burned.items() if t > now} diff --git a/lps/services/search/proxy.py b/lps/services/search/proxy.py index 44e54a9..128d771 100644 --- a/lps/services/search/proxy.py +++ b/lps/services/search/proxy.py @@ -19,10 +19,15 @@ from config.server_configs import decodo_config class DecodoProxy: - def __init__(self, cfg=None, host: str | None = None): + def __init__(self, cfg=None, host: str | None = None, registry=None, owner: str = ""): cfg = cfg if cfg is not None else decodo_config # host 를 넘기면 그 게이트웨이를 쓴다 — 국가 타깃(kr.decodo.com)용. 포트/자격증명은 동일. self.host = host or cfg.host + # registry 를 주면 포트를 **공용 중재자**에게서 배타 임대한다(소스·워커 간 IP 충돌·재사용 방지). + # 없으면 기존 동작(시간창+오프셋 계산) 그대로 — 단독 사용·테스트 경로를 깨지 않는다. + self._registry = registry + self._owner = owner or "proxy" + self._leased: int | None = None self.username = cfg.username self.password = cfg.password self.port_start = cfg.port_start @@ -41,6 +46,9 @@ class DecodoProxy: def rotate(self): """시간창과 무관하게 즉시 다음 포트(=새 IP)로 회전. 봇 감지·예산 도달 시 호출.""" self._rotate_offset += 1 + if self._registry is not None and self._leased is not None: + self._registry.release(self.host, self._leased, self._owner) # 리스 반납 → 다음 _port() 에서 새로 임대 + self._leased = None def seed_offset(self, k: int): """워커별 시작 포트 분산용 — 동시 워커가 같은 포트(=같은 IP)를 쓰지 않도록 시작점을 벌린다.""" @@ -51,11 +59,20 @@ class DecodoProxy: 선제(예산) 회전된 포트는 부르지 않는다 — 불탄 게 아니므로 로테이션 복귀 시 재사용.""" if port is None: return - self._burned[port] = time.monotonic() + (cooldown_sec if cooldown_sec is not None else self.cooldown_sec) - LOG.i(f"[proxy] 포트 {port} 쿨다운 {int(cooldown_sec or self.cooldown_sec)}s — 활성 {self.available_ports()}/{self.port_end - self.port_start + 1}") + cd = cooldown_sec if cooldown_sec is not None else self.cooldown_sec + if self._registry is not None: + # 전역 격리 — 태운 포트를 다른 소스/워커가 곧바로 집는 걸 막는다. + self._registry.burn(self.host, port, cd, owner=self._owner, reason="block") + if self._leased == port: + self._leased = None + return + self._burned[port] = time.monotonic() + cd + LOG.i(f"[proxy] 포트 {port} 쿨다운 {int(cd)}s — 활성 {self.available_ports()}/{self.port_end - self.port_start + 1}") def available_ports(self) -> int: """쿨다운 중이 아닌 포트 수(관측·알림용).""" + if self._registry is not None: + return self._registry.available(self.host) self._prune_burned() return (self.port_end - self.port_start + 1) - len(self._burned) @@ -66,6 +83,18 @@ class DecodoProxy: def _port(self) -> int: """시간창 + 수동 오프셋 기반 포트 선택. 창 안에선 동일 IP, rotate()나 창 변화 시 다음 IP. 쿨다운 중인 포트는 건너뛰고, 전 포트가 쿨다운이면 만료가 가장 임박한 포트를 쓴다(가용성 우선).""" + if self._registry is not None: + lease_sec = self.session_minutes * 60 + if self._leased is not None and self._registry.lease_expired(self.host, self._leased, self._owner, lease_sec): + self._registry.release(self.host, self._leased, self._owner) # sticky 수명 만료 → 새 IP + self._leased = None + if self._leased is None: + # 시작 오프셋은 기존과 동일한 의미(워커별 분산). 실제 중복 방지는 registry 가 보장한다. + bucket = int(time.time() // (self.session_minutes * 60)) + start = (bucket + self._rotate_offset) % (self.port_end - self.port_start + 1) + self._leased = self._registry.acquire(self.host, self._owner, start_at=start, lease_sec=lease_sec) + return self._leased + n = self.port_end - self.port_start + 1 bucket = int(time.time() // (self.session_minutes * 60)) self._prune_burned() diff --git a/lps/tests/test_port_registry.py b/lps/tests/test_port_registry.py new file mode 100644 index 0000000..d85ae87 --- /dev/null +++ b/lps/tests/test_port_registry.py @@ -0,0 +1,160 @@ +"""포트(=IP 세션) 중재자 테스트. + +여기서 지키는 계약은 운영 사고와 직결된다: + - 한 IP 를 두 소유자가 동시에 쓰면 그 IP 가 빨리 탄다 + - 태운 IP 를 다른 소유자가 곧바로 집으면 차단이 전파된다 + - 게이트웨이가 다르면 같은 포트 번호라도 다른 IP 다(실측) — 서로 간섭하면 안 된다 +""" + +import time + +import pytest + +from services.search.port_registry import PortRegistry +from services.search.proxy import DecodoProxy +from config.config_models import DecodoConfig + +GATE, KR = "gate.decodo.com", "kr.decodo.com" + + +def _reg(n: int = 5) -> PortRegistry: + return PortRegistry(10001, 10000 + n) + + +# ── 배타 할당 ──────────────────────────────────────────────────────────── +def test_same_port_is_never_handed_to_two_owners(): + reg = _reg() + taken = {reg.acquire(GATE, f"w{i}") for i in range(5)} + assert len(taken) == 5, "5개 포트를 5명이 나눠 가져야 한다(중복 배정 없음)" + + +def test_released_port_returns_to_pool(): + reg = _reg(2) + a = reg.acquire(GATE, "coupang") + b = reg.acquire(GATE, "naver") + reg.release(GATE, a, "coupang") + assert reg.acquire(GATE, "naver2") == a # 반납된 포트가 다시 나온다 + assert b != a + + +def test_release_by_wrong_owner_is_ignored(): + """뒤늦은 반납이 남의 리스를 깨면 그 순간 두 소유자가 같은 IP 를 쓰게 된다.""" + reg = _reg(1) + p = reg.acquire(GATE, "coupang") + reg.release(GATE, p, "naver") # 남이 반납 시도 → 무시 + assert reg.available(GATE) == 0 + + +# ── 전역 쿨다운 ────────────────────────────────────────────────────────── +def test_burned_port_is_avoided_by_every_owner(): + reg = _reg(2) + p = reg.acquire(GATE, "coupang") + reg.burn(GATE, p, 600, owner="coupang", reason="block") + other = reg.acquire(GATE, "naver") + assert other != p, "다른 소유자도 태운 포트를 피해야 한다" + + +def test_burn_releases_the_lease(): + reg = _reg(3) + p = reg.acquire(GATE, "coupang") + reg.burn(GATE, p, 600) + assert p not in [v for v in reg._held] # 보유 해제 + assert reg.available(GATE) == 2 # 3 - 쿨다운 1 + + +def test_cooldown_expiry_returns_port(): + reg = _reg(1) + p = reg.acquire(GATE, "coupang") + reg.burn(GATE, p, -1) # 이미 만료된 쿨다운 + assert reg.available(GATE) == 1 + assert reg.acquire(GATE, "naver") == p + + +# ── 게이트웨이 격리 ────────────────────────────────────────────────────── +def test_gateways_do_not_interfere(): + """같은 포트 번호라도 gate/kr 은 다른 IP(실측) — 한쪽 점유·차단이 다른 쪽을 막으면 안 된다.""" + reg = _reg(1) + g = reg.acquire(GATE, "coupang") + k = reg.acquire(KR, "naver") + assert g == k == 10001 # 번호는 같아도 서로 다른 자원 + reg.burn(GATE, g, 600) + assert reg.available(GATE) == 0 and reg.available(KR) == 0 # kr 은 여전히 naver 가 보유 + reg.release(KR, k, "naver") + assert reg.available(KR) == 1 + + +# ── 소진 시 동작 ───────────────────────────────────────────────────────── +def test_exhausted_pool_reclaims_soonest_expiring(): + """전부 쿨다운이면 멈추는 대신 가장 빨리 풀릴 포트를 회수한다(가용성 우선).""" + reg = _reg(2) + reg.burn(GATE, 10001, 600) + reg.burn(GATE, 10002, 60) # 얘가 더 빨리 풀린다 + assert reg.acquire(GATE, "naver") == 10002 + + +def test_lease_expiry_is_reported(): + reg = _reg(1) + p = reg.acquire(GATE, "coupang") + assert reg.lease_expired(GATE, p, "coupang", lease_sec=600) is False + assert reg.lease_expired(GATE, p, "naver", lease_sec=600) is True # 소유자 아님 + reg._held[(GATE, p)] = ("coupang", time.monotonic() - 700) + assert reg.lease_expired(GATE, p, "coupang", lease_sec=600) is True + + +def test_snapshot_reports_per_gateway_owners(): + reg = _reg(4) + reg.acquire(GATE, "coupang-w0") + reg.acquire(GATE, "coupang-w1") + reg.acquire(KR, "naver-w0") + reg.burn(KR, 10004, 600) + snap = reg.snapshot() + assert snap[GATE]["held"] == 2 and snap[GATE]["owners"] == {"coupang-w0": 1, "coupang-w1": 1} + assert snap[KR]["held"] == 1 and snap[KR]["cooling"] == 1 + assert snap[KR]["available"] == 2 + + +# ── DecodoProxy 통합 ───────────────────────────────────────────────────── +def _cfg() -> DecodoConfig: + return DecodoConfig(host=GATE, kr_host=KR, username="u", password="p", + port_start=10001, port_end=10003, session_minutes=10) + + +def test_two_sources_never_share_a_port_through_registry(): + reg = _reg(3) + cou = DecodoProxy(_cfg(), registry=reg, owner="coupang-w0") + nav = DecodoProxy(_cfg(), host=KR, registry=reg, owner="naver-w0") + # 게이트웨이가 다르면 번호가 같아도 무방하지만, 같은 게이트웨이를 공유하면 갈려야 한다 + nav_gate = DecodoProxy(_cfg(), registry=reg, owner="naver-w0") + assert cou._port() != nav_gate._port() + assert nav._port() is not None + + +def test_rotate_releases_and_takes_new_port(): + reg = _reg(3) + p = DecodoProxy(_cfg(), registry=reg, owner="coupang-w0") + first = p._port() + p.rotate() + second = p._port() + assert first != second + assert reg.available(GATE) == 2 # 이전 포트는 반납돼 풀에 남는다(쿨다운 아님) + + +def test_mark_burned_goes_global_through_registry(): + reg = _reg(3) + cou = DecodoProxy(_cfg(), registry=reg, owner="coupang-w0") + nav = DecodoProxy(_cfg(), registry=reg, owner="naver-w0") + port = cou._port() + cou.mark_burned(port) + assert cou.available_ports() == 2 # 3 - 쿨다운 1 (보유 없음) + assert nav._port() != port, "쿠팡이 태운 포트를 네이버가 집으면 안 된다" + assert cou.available_ports() == 1 # 3 - 쿨다운 1 - 네이버 보유 1 + + +def test_without_registry_behaviour_is_unchanged(): + """registry 미주입 경로(단독 사용·기존 테스트)는 그대로 동작해야 한다.""" + p = DecodoProxy(_cfg()) + port = p._port() + assert 10001 <= port <= 10003 + p.mark_burned(port) + assert p.available_ports() == 2 + assert p._port() != port diff --git a/lps/worker_main.py b/lps/worker_main.py index 6338480..ce0f29e 100644 --- a/lps/worker_main.py +++ b/lps/worker_main.py @@ -19,6 +19,7 @@ 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 @@ -48,15 +49,18 @@ def _enabled_fallbacks() -> list[str]: return [n for n in names if n in _FALLBACK_SOURCES] -def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history): +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포트를 워커 수로 균등 분할해 시작점을 벌린다. - proxy = DecodoProxy() - # 네이버 전용 프록시 — 같은 계정·같은 포트지만 **한국 타깃 게이트웨이**를 쓴다. - # 해외 residential IP 로는 msearch 가 즉시 하드차단된다(2026-08-04 실측). kr_host 가 비면 - # 기본 게이트웨이로 떨어지고, 그때는 네이버가 막힐 수 있다(로그의 BOT-DETECTED 로 드러남). - kr_proxy = DecodoProxy(host=decodo_config.kr_host or None) + # 두 프록시는 **게이트웨이가 다르다** — 쿠팡=국가 무지정, 네이버=한국 전용. + # 해외 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)) @@ -123,8 +127,10 @@ async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float def _proxy_ports_snapshot(adapters) -> tuple[int, int] | None: - """워커 프록시들의 가용 포트 현황 — (최소 가용 수, 전체 포트 수). 프록시 미사용이면 None. - 쿨다운 맵은 프록시 인스턴스(워커)별이라 가장 소진된 워커 기준(min)으로 본다.""" + """가용 포트 현황 — (최소 가용 수, 전체 포트 수). 프록시 미사용이면 None. + + 게이트웨이가 둘(쿠팡=국가무지정 / 네이버=한국)이라 **가장 마른 게이트웨이 기준(min)**으로 본다 + — 한쪽만 고갈돼도 그 소스는 검색을 못 하므로 평균으로 덮으면 안 된다.""" proxies = {} for ad in (adapters or []): p = getattr(ad, "_proxy", None) @@ -137,6 +143,16 @@ def _proxy_ports_snapshot(adapters) -> tuple[int, int] | None: 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 로 발화 @@ -171,8 +187,12 @@ async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0, 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} — 대규모 차단 진행 신호", snap) + f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호" + + (f" · {pool}" if pool else ""), snap) # 예산 누수 — 요청 예산을 지켰는데도 차단된 IP 세션 발생. # 처방은 '몇 번째 요청에서 막혔나'로 갈린다: # ip_req#1 위주 → 새 IP 첫 요청부터 차단 = IP 평판 문제. 예산을 낮춰도 소용없다. @@ -270,8 +290,10 @@ async def main(concurrency: int = 1): # 잡 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) + 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()