diff --git a/lps/common/database/model/models.py b/lps/common/database/model/models.py index 2e96c5e..5f440d1 100644 --- a/lps/common/database/model/models.py +++ b/lps/common/database/model/models.py @@ -152,3 +152,42 @@ class bot_detection(MAIN_BASE): headless = Column(Boolean, nullable=True) html_len = Column(Integer, nullable=True) # 응답 길이(차단 페이지는 작음) created_at = Column(DateTime(timezone=True), nullable=False, server_default=text("now()")) + + +class proxy_port(MAIN_BASE): + """프록시 포트(=IP 세션) 임대 장부 — **프로세스 간 공유 상태**. + + 한 DECODO 계정을 여러 프로세스(워커 컨테이너·PROCESS_COUNT)가 나눠 쓰기 때문에 + 임대·쿨다운·휴식을 프로세스 메모리에 두면 서로의 상태를 모른다. 같은 IP 를 동시에 + 잡거나, 한쪽이 태운 IP 를 다른 쪽이 곧바로 집는다. 그래서 DB 를 단일 진실로 둔다. + + 행 1개 = 게이트웨이의 포트 1개(=sticky IP 세션 1개). 상태는 세 시각으로만 표현한다: + leased_until 임대 중(만료되면 자동 해제 — 프로세스가 죽어도 IP 가 영구히 묶이지 않는다) + rest_until 선제 회전으로 쉬는 중(탄 게 아님) + cooldown_until 차단당해 격리 중 + 셋 다 지났으면 가용. 회전은 last_used_at 오래된 순(LRU)이라 프로세스가 늘어도 + 전체가 자연스럽게 한 바퀴씩 돈다. + """ + + @staticmethod + def DBType(): + return DBType.MAIN.value + + __tablename__ = "proxy_port" + + host = Column(String(80), primary_key=True) # 게이트웨이(gate/kr — 같은 번호라도 IP 가 다름) + port = Column(Integer, primary_key=True) + owner = Column(String(80), nullable=True) # 현재 임대자(worker/소스 식별) + leased_until = Column(DateTime(timezone=True), nullable=True) # 임대 만료(=sticky 수명) + rest_until = Column(DateTime(timezone=True), nullable=True) # 휴식 만료(선제 회전) + cooldown_until = Column(DateTime(timezone=True), nullable=True) # 쿨다운 만료(차단) + last_used_at = Column(DateTime(timezone=True), nullable=True) # 마지막 임대 시각(LRU 회전 기준) + last_reason = Column(String(40), nullable=True) # 마지막 상태 변경 사유(block/budget/window…) + use_count = Column(Integer, nullable=False, server_default=text("0")) # 누적 임대 횟수(관측) + burn_count = Column(Integer, nullable=False, server_default=text("0")) # 누적 차단 횟수(불량 IP 슬롯 식별) + updated_at = Column(DateTime(timezone=True), nullable=False, server_default=text("now()"), onupdate=text("now()")) + + __table_args__ = ( + # acquire 정렬/필터용 — 가용 판정(3개 시각)과 LRU 정렬을 한 인덱스로 태운다. + Index("ix_proxy_port_pick", "host", "last_used_at"), + ) diff --git a/lps/crud/port_lease.py b/lps/crud/port_lease.py new file mode 100644 index 0000000..0a8b730 --- /dev/null +++ b/lps/crud/port_lease.py @@ -0,0 +1,146 @@ +"""프록시 포트(=IP 세션) 임대 — **프로세스 간 공유** 장부(DB 단일 진실). + +한 DECODO 계정을 여러 프로세스가 나눠 쓴다. 인메모리 장부로는 서로의 임대·차단을 모르므로 +같은 IP 를 동시에 잡거나 태운 IP 를 곧바로 재사용한다. 그래서 큐(job)와 같은 방식으로 푼다: +**한 UPDATE 안에서 FOR UPDATE SKIP LOCKED 로 후보를 잠그고 임대까지 끝낸다**(원자적, 이중 배정 불가). + +회전 정책은 LRU(last_used_at 오래된 순)다. 프로세스가 몇 개든 각자 '가장 오래 안 쓴 IP'를 +집으므로 전체가 자연히 한 바퀴씩 돈다 — 프로세스별 오프셋 계산이 필요 없다. + +죽은 프로세스 회수: 임대는 leased_until 로 만료된다. 워커가 죽어도 그 IP 는 sticky 수명이 +지나면 자동으로 풀린다(잡 큐의 lease/reaper 와 같은 발상 — 별도 정리 프로세스가 필요 없다). +""" + +from sqlalchemy import text + +from common.database.db_session_manager import DB_SESSION_MNG +from common.enums import DBType, DBWRType + +# 가용 판정 — 세 시각이 모두 지났으면 쓸 수 있다. +_FREE = """ + (leased_until IS NULL OR leased_until < now()) +AND (rest_until IS NULL OR rest_until < now()) +AND (cooldown_until IS NULL OR cooldown_until < now()) +""" + + +class PortLeaseStore: + DB = DBType.MAIN.value + + async def _tx(self, fn): + s = await DB_SESSION_MNG.start_session(self.DB, DBWRType.DB_WRITE.value) + try: + out = await fn(s) + await s.commit() + return out + except Exception: + await s.rollback() + raise + finally: + await DB_SESSION_MNG.end_session(self.DB, DBWRType.DB_WRITE.value) + + async def ensure_ports(self, host: str, port_start: int, port_end: int): + """게이트웨이의 포트 행을 보장(최초 1회). 이미 있으면 그대로 둔다 — 상태를 덮으면 안 된다.""" + sql = text(""" + INSERT INTO proxy_port (host, port) + SELECT :host, g FROM generate_series(CAST(:s AS int), CAST(:e AS int)) AS g + ON CONFLICT (host, port) DO NOTHING + """) + await self._tx(lambda s: s.execute(sql, {"host": host, "s": port_start, "e": port_end})) + + async def acquire(self, host: str, owner: str, lease_sec: float) -> int | None: + """가용 포트 1개를 원자적으로 임대(LRU). 전부 막혀 있으면 None. + + 후보 선택과 임대를 한 문장에서 끝낸다 — SELECT 후 UPDATE 로 나누면 그 사이에 + 다른 프로세스가 같은 행을 집어 같은 IP 를 동시에 쓰게 된다. + """ + sql = text(f""" + UPDATE proxy_port p SET + owner = :owner, + leased_until = now() + make_interval(secs => CAST(:lease AS double precision)), + last_used_at = now(), + rest_until = NULL, + use_count = p.use_count + 1, + last_reason = 'acquire', + updated_at = now() + WHERE (p.host, p.port) = ( + SELECT c.host, c.port FROM proxy_port c + WHERE c.host = :host AND {_FREE} + ORDER BY c.last_used_at NULLS FIRST, c.port + FOR UPDATE SKIP LOCKED + LIMIT 1 + ) + RETURNING p.port + """) + + async def run(s): + row = (await s.execute(sql, {"host": host, "owner": owner, "lease": lease_sec})).first() + return row[0] if row else None + + return await self._tx(run) + + async def release(self, host: str, port: int, owner: str, rest_sec: float = 0): + """임대 반납. rest_sec>0 이면 그만큼 휴식(선제 회전) — 다른 프로세스도 그동안 못 집는다. + 소유자가 다르면 무시한다(뒤늦은 반납이 남의 임대를 깨지 않도록).""" + sql = text(""" + UPDATE proxy_port SET + owner = NULL, + leased_until = NULL, + rest_until = CASE WHEN CAST(:rest AS double precision) > 0 THEN now() + make_interval(secs => CAST(:rest AS double precision)) ELSE NULL END, + last_reason = CASE WHEN CAST(:rest AS double precision) > 0 THEN 'rest' ELSE 'release' END, + updated_at = now() + WHERE host = :host AND port = :port AND owner = :owner + """) + await self._tx(lambda s: s.execute(sql, {"host": host, "port": port, "owner": owner, "rest": rest_sec})) + + async def burn(self, host: str, port: int, cooldown_sec: float, owner: str = "", reason: str = "block"): + """차단된 포트를 전역 격리. 소유자와 무관하게 적용한다 — 차단은 사실이지 소유권 문제가 아니다.""" + sql = text(""" + UPDATE proxy_port SET + owner = NULL, + leased_until = NULL, + rest_until = NULL, + cooldown_until = now() + make_interval(secs => CAST(:cd AS double precision)), + burn_count = burn_count + 1, + last_reason = :reason, + updated_at = now() + WHERE host = :host AND port = :port + """) + await self._tx(lambda s: s.execute(sql, {"host": host, "port": port, "cd": cooldown_sec, "reason": reason})) + + async def renew(self, host: str, port: int, owner: str, lease_sec: float) -> bool: + """임대 연장. 아직 내 것이면 True — False 면 만료·회수된 것이므로 새로 잡아야 한다.""" + sql = text(""" + UPDATE proxy_port SET leased_until = now() + make_interval(secs => CAST(:lease AS double precision)), updated_at = now() + WHERE host = :host AND port = :port AND owner = :owner + AND (leased_until IS NULL OR leased_until > now()) + """) + + async def run(s): + return (await s.execute(sql, {"host": host, "port": port, "owner": owner, "lease": lease_sec})).rowcount + + return bool(await self._tx(run)) + + async def snapshot(self, host: str | None = None) -> dict: + """게이트웨이별 현황 — ops 알림/관리자 화면용.""" + sql = text(f""" + SELECT host, + count(*) FILTER (WHERE leased_until > now()) AS held, + count(*) FILTER (WHERE rest_until > now()) AS resting, + count(*) FILTER (WHERE cooldown_until > now()) AS cooling, + count(*) FILTER (WHERE {_FREE}) AS available, + count(*) AS total + FROM proxy_port + WHERE (CAST(:host AS varchar) IS NULL OR host = :host) + GROUP BY host ORDER BY host + """) + + async def run(s): + rows = (await s.execute(sql, {"host": host})).mappings().all() + return {r["host"]: {k: r[k] for k in ("held", "resting", "cooling", "available", "total")} for r in rows} + + return await self._tx(run) + + async def available(self, host: str) -> int: + snap = await self.snapshot(host) + return snap.get(host, {}).get("available", 0) diff --git a/lps/migrations/2026-08-05-proxy_port.sql b/lps/migrations/2026-08-05-proxy_port.sql new file mode 100644 index 0000000..4d7b27d --- /dev/null +++ b/lps/migrations/2026-08-05-proxy_port.sql @@ -0,0 +1,18 @@ +-- 프록시 포트(=IP 세션) 임대 장부 — 프로세스 간 공유 상태. +-- 한 DECODO 계정을 여러 워커 프로세스가 나눠 쓰므로 임대·휴식·쿨다운을 DB 에 둔다. +-- (인메모리면 서로의 차단을 몰라 같은 IP 를 동시에 잡거나 태운 IP 를 곧바로 재사용한다) +CREATE TABLE IF NOT EXISTS proxy_port ( + host varchar(80) NOT NULL, -- 게이트웨이(gate/kr — 같은 번호라도 IP 가 다름) + port integer NOT NULL, + owner varchar(80), -- 현재 임대자(소스-PID-워커) + leased_until timestamptz, -- 임대 만료(=sticky 수명). 프로세스가 죽어도 자동 회수 + rest_until timestamptz, -- 휴식 만료(선제 회전 — 탄 게 아님) + cooldown_until timestamptz, -- 쿨다운 만료(차단) + last_used_at timestamptz, -- 마지막 임대 시각(LRU 회전 기준) + last_reason varchar(40), + use_count integer NOT NULL DEFAULT 0, + burn_count integer NOT NULL DEFAULT 0, -- 누적 차단(불량 IP 슬롯 식별) + updated_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (host, port) +); +CREATE INDEX IF NOT EXISTS ix_proxy_port_pick ON proxy_port (host, last_used_at); diff --git a/lps/services/search/browser_base.py b/lps/services/search/browser_base.py index 26422e0..123fa3e 100644 --- a/lps/services/search/browser_base.py +++ b/lps/services/search/browser_base.py @@ -170,6 +170,10 @@ class BrowserSearchAdapter(SearchAdapter): else: kwargs["channel"] = _CHROME_CHANNEL # 로컬: 실제 Chrome if self._proxy and self._proxy.enabled: + # 포트 확정은 여기서만 한다(유일한 async 지점) — 회전·차단으로 밀린 DB 반영도 함께 flush. + ensure = getattr(self._proxy, "ensure_port", None) + if ensure is not None: + await ensure() kwargs["proxy"] = self._proxy.playwright_proxy() self._current_port = self._proxy.current_port self._ctx = await self._pw.chromium.launch_persistent_context(**kwargs) @@ -182,6 +186,14 @@ class BrowserSearchAdapter(SearchAdapter): self._force_recycle = False async def _close_ctx(self): + # 밀린 반납·차단을 먼저 DB 에 반영한다 — 여기서 흘리지 않으면 종료 시 유실돼 + # 다른 프로세스가 방금 태운 IP 를 그대로 집는다. + flush = getattr(self._proxy, "flush", None) + if flush is not None: + try: + await flush() + except Exception as ex: + LOG.d(f"[{self.source}] 포트 상태 flush 실패(무시): {type(ex).__name__}") self._cdp = None # 컨텍스트와 함께 CDP 세션도 죽음 → 다음 검색 때 재부착 if self._ctx is not None: await self._record_session_end() diff --git a/lps/services/search/profile_slot.py b/lps/services/search/profile_slot.py new file mode 100644 index 0000000..d55a0e8 --- /dev/null +++ b/lps/services/search/profile_slot.py @@ -0,0 +1,46 @@ +"""Chrome 프로필 슬롯 배정 — 프로세스가 여러 개여도 같은 프로필을 잡지 않게. + +Chrome 은 user_data_dir 당 인스턴스 하나만 허용한다(ProcessSingleton). 워커 프로세스를 +여러 개 띄우면 두 번째가 "기존 브라우저 세션에서 여는 중입니다" 로 죽고, 그 소스의 검색이 +전부 TargetClosedError 로 실패한다(2026-08-05 실측: 워커 2프로세스 → 잡 2건 DEAD). + +PID 를 경로에 넣는 방식은 쓰지 않는다 — 재시작마다 새 디렉터리가 생겨 **웜 쿠키 +(cf_clearance·Akamai 세션)를 매번 버리고** 디렉터리가 무한히 쌓인다. 대신 슬롯 번호를 +파일 락으로 선점한다: 살아있는 프로세스가 없으면 같은 슬롯을 다시 쓰므로 쿠키가 유지되고, +프로세스가 죽으면 OS 가 락을 자동 해제해 슬롯이 곧바로 회수된다(정리 코드 불필요). +""" + +import fcntl +import os + +from common.logger import LOG + +# 잡은 락 핸들을 프로세스 수명 동안 살려 둔다 — GC 되면 락이 풀려 다른 프로세스가 같은 슬롯을 집는다. +_HELD: list = [] + + +def claim_profile_slot(base_dir: str, source: str, worker_index: int = 0, max_slots: int = 32) -> str: + """이 프로세스가 배타적으로 쓸 프로필 경로를 돌려준다. + + 같은 (source, worker_index) 라도 프로세스가 다르면 다른 슬롯이 배정된다. + 슬롯을 모두 뺏겼으면 마지막 후보를 그대로 쓴다 — 검색을 못 하는 것보다 낫고, + 그 경우엔 로그로 드러난다(동시 프로세스 수가 max_slots 를 넘었다는 뜻). + """ + os.makedirs(base_dir, exist_ok=True) + for slot in range(max_slots): + path = os.path.join(base_dir, f"lps_{source}_w{worker_index}_s{slot}") + os.makedirs(path, exist_ok=True) + lock_path = path + ".lock" + try: + fh = open(lock_path, "w") + fcntl.flock(fh, fcntl.LOCK_EX | fcntl.LOCK_NB) # 비블로킹 — 남이 쥐고 있으면 즉시 실패 + except OSError: + continue + fh.write(f"{os.getpid()}\n") + fh.flush() + _HELD.append(fh) + return path + + fallback = os.path.join(base_dir, f"lps_{source}_w{worker_index}_s{max_slots - 1}") + LOG.w(f"[profile] {source} 슬롯 {max_slots}개 모두 사용 중 — {fallback} 공유(프로세스 과다)") + return fallback diff --git a/lps/services/search/proxy.py b/lps/services/search/proxy.py index e18b306..a63597b 100644 --- a/lps/services/search/proxy.py +++ b/lps/services/search/proxy.py @@ -19,15 +19,23 @@ from config.server_configs import decodo_config class DecodoProxy: - def __init__(self, cfg=None, host: str | None = None, registry=None, owner: str = ""): + def __init__(self, cfg=None, host: str | None = None, registry=None, owner: str = "", store=None): cfg = cfg if cfg is not None else decodo_config # host 를 넘기면 그 게이트웨이를 쓴다 — 국가 타깃(kr.decodo.com)용. 포트/자격증명은 동일. self.host = host or cfg.host - # registry 를 주면 포트를 **공용 중재자**에게서 배타 임대한다(소스·워커 간 IP 충돌·재사용 방지). - # 없으면 기존 동작(시간창+오프셋 계산) 그대로 — 단독 사용·테스트 경로를 깨지 않는다. + # 포트 배정 방식 3가지(우선순위 순): + # store DB 장부(PortLeaseStore) — **프로세스 간 공유**. 여러 워커 프로세스가 한 계정을 + # 나눠 쓸 때 유일하게 안전하다. 프로덕션 경로. + # registry 인메모리 중재자 — 단일 프로세스 전용(테스트·단독 실행) + # 둘 다 없으면 기존 계산식(시간창+오프셋) + self._store = store self._registry = registry self._owner = owner or "proxy" self._leased: int | None = None + # DB 왕복은 비동기라 검색 루프(동기 호출부)에서 곧바로 못 한다. 회전·차단은 여기에 적어두고 + # 다음 ensure_port(브라우저 재기동 직전, async)에서 한 번에 반영한다. + self._pending_release: tuple[int, float] | None = None # (port, rest_sec) + self._pending_burn: list[tuple[int, float]] = [] # [(port, cooldown_sec)] self.username = cfg.username self.password = cfg.password self.port_start = cfg.port_start @@ -54,8 +62,13 @@ class DecodoProxy: 차단으로 태우는 건 mark_burned 가 따로 처리한다(휴식보다 훨씬 긴 쿨다운). """ self._rotate_offset += 1 - if self._registry is not None and self._leased is not None: - rest = self.rest_sec if kind == "budget" else 0 + if self._leased is None: + return + rest = self.rest_sec if kind == "budget" else 0 + if self._store is not None: + self._pending_release = (self._leased, rest) # 다음 ensure_port 에서 DB 반영 + self._leased = None + elif self._registry is not None: self._registry.release(self.host, self._leased, self._owner, rest_sec=rest) self._leased = None @@ -69,6 +82,12 @@ class DecodoProxy: if port is None: return cd = cooldown_sec if cooldown_sec is not None else self.cooldown_sec + if self._store is not None: + self._pending_burn.append((port, cd)) + if self._leased == port: + self._leased = None + LOG.i(f"[port] {self.host}:{port} 차단 기록 예약 {int(cd)}s ({self._owner})") + return if self._registry is not None: # 전역 격리 — 태운 포트를 다른 소스/워커가 곧바로 집는 걸 막는다. self._registry.burn(self.host, port, cd, owner=self._owner, reason="block") @@ -78,6 +97,45 @@ class DecodoProxy: self._burned[port] = time.monotonic() + cd LOG.i(f"[proxy] 포트 {port} 쿨다운 {int(cd)}s — 활성 {self.available_ports()}/{self.port_end - self.port_start + 1}") + async def ensure_port(self) -> int | None: + """다음 요청에 쓸 포트를 확정한다(브라우저 재기동 직전에 await). + + DB 장부 모드에서 이 함수가 유일한 I/O 지점이다: 밀린 반납·차단을 먼저 flush 하고, + 임대가 없으면 새로 잡는다. 임대가 살아 있으면 연장(renew)해 sticky 수명 동안 붙잡는다 + — 연장에 실패하면 남이 회수해 간 것이므로 새 포트를 잡는다. + """ + if self._store is None: + return self._port() if self.enabled else None + + if self._pending_release is not None: # 선제 회전/일반 반납 + port, rest = self._pending_release + self._pending_release = None + await self._store.release(self.host, port, self._owner, rest_sec=rest) + while self._pending_burn: # 차단 격리(전역) + port, cd = self._pending_burn.pop(0) + await self._store.burn(self.host, port, cd, owner=self._owner, reason="block") + + lease_sec = self.session_minutes * 60 + if self._leased is not None and not await self._store.renew(self.host, self._leased, self._owner, lease_sec): + self._leased = None # 만료·회수됨 → 새로 잡는다 + if self._leased is None: + self._leased = await self._store.acquire(self.host, self._owner, lease_sec) + if self._leased is None: + LOG.w(f"[port] {self.host} 가용 포트 없음 — 전부 임대/휴식/쿨다운 중({self._owner})") + return self._leased + + async def flush(self): + """밀린 반납·차단만 반영(검색 종료·셧다운 시). 포트를 새로 잡지는 않는다.""" + if self._store is None: + return + if self._pending_release is not None: + port, rest = self._pending_release + self._pending_release = None + await self._store.release(self.host, port, self._owner, rest_sec=rest) + while self._pending_burn: + port, cd = self._pending_burn.pop(0) + await self._store.burn(self.host, port, cd, owner=self._owner, reason="block") + def available_ports(self) -> int: """쿨다운 중이 아닌 포트 수(관측·알림용).""" if self._registry is not None: @@ -92,6 +150,11 @@ class DecodoProxy: def _port(self) -> int: """시간창 + 수동 오프셋 기반 포트 선택. 창 안에선 동일 IP, rotate()나 창 변화 시 다음 IP. 쿨다운 중인 포트는 건너뛰고, 전 포트가 쿨다운이면 만료가 가장 임박한 포트를 쓴다(가용성 우선).""" + if self._store is not None: + # 임대는 ensure_port(async)가 확정한다. 여기선 확정값을 돌려줄 뿐 — 동기 경로에서 + # DB 를 만지지 않는다. 아직 못 잡았으면 계산식으로 폴백(로그·프리플라이트용). + if self._leased is not None: + return self._leased 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): diff --git a/lps/tests/test_alerts.py b/lps/tests/test_alerts.py index 7ab0b9d..2850eb1 100644 --- a/lps/tests/test_alerts.py +++ b/lps/tests/test_alerts.py @@ -101,27 +101,43 @@ def test_pool_status_shape(): # ---- 가용 프록시 포트 현황(포트 고갈 룰의 데이터) -------------------------- -def test_proxy_ports_snapshot_min_across_workers(): - from config.config_models import DecodoConfig - from services.search.proxy import DecodoProxy - from worker_main import _proxy_ports_snapshot +async def test_port_pool_status_picks_worst_gateway(): + """게이트웨이가 둘이라 **가장 마른 쪽** 기준으로 알림을 걸어야 한다 — + 한쪽만 고갈돼도 그 소스(네이버/쿠팡)는 검색을 못 하므로 평균으로 덮으면 안 된다.""" + from worker_main import _port_pool_status - def _proxy(): - return DecodoProxy(DecodoConfig(host="h", username="u", password="p", - port_start=10001, port_end=10010, session_minutes=10)) + class _Store: + async def snapshot(self, host=None): + return {"gate.decodo.com": {"held": 2, "resting": 0, "cooling": 0, "available": 8, "total": 10}, + "kr.decodo.com": {"held": 1, "resting": 2, "cooling": 4, "available": 3, "total": 10}} + + class _Proxy: + _store = _Store() class _Ad: - def __init__(self, proxy): - self._proxy = proxy + _proxy = _Proxy() - p1, p2 = _proxy(), _proxy() - p2.mark_burned(10001) - p2.mark_burned(10002) - avail, total = _proxy_ports_snapshot([_Ad(p1), _Ad(p2), object()]) # 프록시 없는 어댑터 혼재 OK - assert (avail, total) == (8, 10) # 가장 소진된 워커(p2) 기준 + detail, ports = await _port_pool_status([object(), _Ad()]) # 프록시 없는 어댑터 혼재 OK + assert ports == (3, 10) # kr(=3) 기준 + assert detail["kr.decodo.com"]["cooling"] == 4 # 원인 파악용 상세도 함께 실린다 -def test_proxy_ports_snapshot_none_without_proxy(): - from worker_main import _proxy_ports_snapshot - assert _proxy_ports_snapshot([object()]) is None - assert _proxy_ports_snapshot(None) is None +async def test_port_pool_status_none_without_store(): + from worker_main import _port_pool_status + assert await _port_pool_status([object()]) == ({}, None) + assert await _port_pool_status(None) == ({}, None) + + +async def test_port_pool_status_survives_db_error(): + """장부 조회가 실패해도 ops 모니터는 계속 돌아야 한다(알림만 건너뛴다).""" + from worker_main import _port_pool_status + + class _Broken: + async def snapshot(self, host=None): + raise RuntimeError("db down") + + class _Ad: + class _proxy: + _store = _Broken() + + assert await _port_pool_status([_Ad()]) == ({}, None) diff --git a/lps/tests/test_port_lease.py b/lps/tests/test_port_lease.py new file mode 100644 index 0000000..096e2ca --- /dev/null +++ b/lps/tests/test_port_lease.py @@ -0,0 +1,138 @@ +"""프로세스 간 공유 포트 장부(proxy_port) 테스트 — 실제 DB 필요. + +여기서 지키는 계약은 멀티 프로세스 운영과 직결된다: + - 두 프로세스가 같은 IP 를 동시에 잡으면 그 IP 가 두 배 속도로 탄다 + - 한 프로세스가 태운 IP 를 다른 프로세스가 곧바로 집으면 차단이 전파된다 + - 워커가 죽은 채로 임대를 쥐고 있으면 그 IP 는 영영 안 돌아온다(임대 만료로 자동 회수돼야 함) +""" + +import pytest +from sqlalchemy import text + +from common.database.db_session_manager import DB_SESSION_MNG +from common.enums import DBType, DBWRType +from crud.port_lease import PortLeaseStore + +HOST_A, HOST_B = "test-gate.example", "test-kr.example" +P_START, P_END = 20001, 20005 + + +@pytest.fixture +async def store(db_engine): + st = PortLeaseStore() + await _exec("DELETE FROM proxy_port WHERE host IN (:a, :b)", {"a": HOST_A, "b": HOST_B}) + await st.ensure_ports(HOST_A, P_START, P_END) + await st.ensure_ports(HOST_B, P_START, P_END) + yield st + await _exec("DELETE FROM proxy_port WHERE host IN (:a, :b)", {"a": HOST_A, "b": HOST_B}) + + +async def _exec(sql: str, params: dict | None = None): + s = await DB_SESSION_MNG.start_session(DBType.MAIN.value, DBWRType.DB_WRITE.value) + try: + await s.execute(text(sql), params or {}) + await s.commit() + finally: + await DB_SESSION_MNG.end_session(DBType.MAIN.value, DBWRType.DB_WRITE.value) + + +# ── 배타 임대(프로세스 간) ─────────────────────────────────────────────── +async def test_no_two_owners_get_the_same_port(store): + """서로 다른 프로세스를 흉내낸 owner 5개가 5포트를 하나씩 나눠 가져야 한다.""" + got = [await store.acquire(HOST_A, f"p{i}-w0", 600) for i in range(5)] + assert None not in got + assert len(set(got)) == 5 + + +async def test_pool_exhaustion_returns_none(store): + for i in range(5): + assert await store.acquire(HOST_A, f"p{i}", 600) is not None + assert await store.acquire(HOST_A, "p-late", 600) is None # 남는 포트 없음 → 호출부가 대기/재시도 + + +async def test_gateways_are_independent(store): + """같은 포트 번호라도 게이트웨이가 다르면 다른 IP(실측) — 한쪽 임대가 다른 쪽을 막으면 안 된다.""" + a = await store.acquire(HOST_A, "coupang", 600) + b = await store.acquire(HOST_B, "naver", 600) + assert a == b == P_START # 번호는 같아도 서로 다른 자원 + assert (await store.available(HOST_A)) == 4 and (await store.available(HOST_B)) == 4 + + +# ── 차단·휴식(전역) ────────────────────────────────────────────────────── +async def test_burned_port_is_invisible_to_every_process(store): + p = await store.acquire(HOST_A, "coupang-p1", 600) + await store.burn(HOST_A, p, 1800, owner="coupang-p1") + others = {await store.acquire(HOST_A, f"naver-p{i}", 600) for i in range(4)} + assert p not in others # 다른 프로세스도 태운 IP 를 못 집는다 + assert await store.acquire(HOST_A, "naver-p9", 600) is None + + +async def test_rested_port_is_held_back_then_returns(store): + p = await store.acquire(HOST_A, "naver-p1", 600) + await store.release(HOST_A, p, "naver-p1", rest_sec=600) + others = {await store.acquire(HOST_A, f"x{i}", 600) for i in range(4)} + assert p not in others # 쉬는 동안은 아무도 못 집는다 + await store.release(HOST_A, p, "nobody", rest_sec=0) # 소유자 불일치 → 무시돼야 함 + assert await store.acquire(HOST_A, "y", 600) is None + await _exec("UPDATE proxy_port SET rest_until = now() - interval '1 s' WHERE host=:h AND port=:p", + {"h": HOST_A, "p": p}) + assert await store.acquire(HOST_A, "z", 600) == p # 휴식 만료 → 복귀 + + +async def test_release_without_rest_returns_port_to_pool(store): + """반납된 포트는 곧바로 풀에 돌아온다. 단 **다음 차례가 되는 건 아니다** — + LRU 라 한 번도 안 쓴 포트가 먼저 나가고, 반납분은 한 바퀴 뒤에 다시 온다(회전의 정의).""" + p = await store.acquire(HOST_A, "naver-p1", 600) + await store.release(HOST_A, p, "naver-p1") + assert await store.available(HOST_A) == 5 # 5개 전부 가용 + + picked = [await store.acquire(HOST_A, f"w{i}", 600) for i in range(5)] + assert p in picked and picked[-1] == p # 가장 최근 사용분이 맨 마지막 + + +# ── 죽은 프로세스 회수 ─────────────────────────────────────────────────── +async def test_expired_lease_is_reclaimed_without_a_reaper(store): + """워커가 임대를 쥔 채 죽어도 sticky 수명이 지나면 자동으로 풀려야 한다.""" + p = await store.acquire(HOST_A, "dead-process", 600) + assert await store.available(HOST_A) == 4 # 죽은 프로세스가 쥔 동안은 빠져 있고 + await _exec("UPDATE proxy_port SET leased_until = now() - interval '1 s' WHERE host=:h AND port=:p", + {"h": HOST_A, "p": p}) + assert await store.available(HOST_A) == 5 # 임대 만료 → 별도 정리 없이 자동 복귀 + picked = [await store.acquire(HOST_A, f"alive{i}", 600) for i in range(5)] + assert p in picked + + +async def test_renew_keeps_the_lease_and_fails_after_takeover(store): + p = await store.acquire(HOST_A, "p1", 600) + assert await store.renew(HOST_A, p, "p1", 600) is True + await _exec("UPDATE proxy_port SET owner='p2' WHERE host=:h AND port=:p", {"h": HOST_A, "p": p}) + assert await store.renew(HOST_A, p, "p1", 600) is False # 남이 가져감 → 새로 잡아야 함 + + +# ── 회전(LRU) ──────────────────────────────────────────────────────────── +async def test_rotation_is_least_recently_used(store): + """프로세스가 몇 개든 '가장 오래 안 쓴 IP'를 집으므로 전체가 한 바퀴씩 돈다.""" + first = [] + for i in range(5): # 5포트를 한 바퀴 소진 + p = await store.acquire(HOST_A, f"w{i}", 600) + first.append(p) + await store.release(HOST_A, p, f"w{i}") + second = [] + for i in range(5): + p = await store.acquire(HOST_A, f"w{i}", 600) + second.append(p) + await store.release(HOST_A, p, f"w{i}") + assert sorted(first) == sorted(second) # 같은 5개를 + assert first == second # 같은 순서로(LRU) — 특정 IP 편중 없음 + + +async def test_snapshot_counts_each_state(store): + held = await store.acquire(HOST_A, "w0", 600) + rested = await store.acquire(HOST_A, "w1", 600) + await store.release(HOST_A, rested, "w1", rest_sec=600) + burned = await store.acquire(HOST_A, "w2", 600) + await store.burn(HOST_A, burned, 1800) + snap = (await store.snapshot(HOST_A))[HOST_A] + assert snap["held"] == 1 and snap["resting"] == 1 and snap["cooling"] == 1 + assert snap["available"] == 2 and snap["total"] == 5 + assert held != rested != burned diff --git a/lps/tests/test_profile_slot.py b/lps/tests/test_profile_slot.py new file mode 100644 index 0000000..daadf6f --- /dev/null +++ b/lps/tests/test_profile_slot.py @@ -0,0 +1,34 @@ +"""Chrome 프로필 슬롯 배정 — 프로세스가 여러 개여도 같은 프로필을 잡으면 안 된다.""" + +import os + +from services.search import profile_slot +from services.search.profile_slot import claim_profile_slot + + +def test_second_claim_gets_a_different_slot(tmp_path): + """같은 소스·같은 워커 인덱스라도 두 번째 요청(=다른 프로세스)은 다른 슬롯을 받아야 한다.""" + a = claim_profile_slot(str(tmp_path), "coupang", worker_index=0) + b = claim_profile_slot(str(tmp_path), "coupang", worker_index=0) + assert a != b and os.path.isdir(a) and os.path.isdir(b) + + +def test_sources_and_workers_are_separate(tmp_path): + c = claim_profile_slot(str(tmp_path), "coupang", worker_index=0) + n = claim_profile_slot(str(tmp_path), "naver", worker_index=0) + w1 = claim_profile_slot(str(tmp_path), "coupang", worker_index=1) + assert len({c, n, w1}) == 3 + + +def test_slot_is_reused_after_release(tmp_path): + """프로세스가 죽으면 OS 가 락을 풀고 같은 슬롯이 재사용된다 — 웜 쿠키를 버리지 않기 위함.""" + a = claim_profile_slot(str(tmp_path), "coupang", worker_index=0) + for fh in list(profile_slot._HELD): # 프로세스 종료를 흉내 + fh.close() + profile_slot._HELD.remove(fh) + assert claim_profile_slot(str(tmp_path), "coupang", worker_index=0) == a + + +def test_exhausted_slots_fall_back_without_raising(tmp_path): + paths = [claim_profile_slot(str(tmp_path), "coupang", worker_index=0, max_slots=2) for _ in range(3)] + assert paths[2] == paths[1] # 마지막 슬롯 공유(검색 중단보다 낫다) diff --git a/lps/worker_main.py b/lps/worker_main.py index d2333e2..652d38f 100644 --- a/lps/worker_main.py +++ b/lps/worker_main.py @@ -18,8 +18,9 @@ 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.port_lease import PortLeaseStore from crud.price_history import PriceHistory -from services.search.port_registry import PortRegistry +from services.search.profile_slot import claim_profile_slot from services.search.proxy import DecodoProxy from services.search.coupang.adapter import CoupangAdapter from services.search.naver_shop.adapter import NaverShopAdapter @@ -49,7 +50,7 @@ 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, registry=None): +def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history, store=None): """워커 1개의 자립 세트(브라우저 어댑터·AI·핸들러)를 만든다. 프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP)로 진짜 병렬을 보장한다.""" # 워커별 프록시(다른 포트=다른 IP). 100포트를 워커 수로 균등 분할해 시작점을 벌린다. @@ -57,21 +58,25 @@ def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history # 해외 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)) + # 포트는 **DB 장부(proxy_port)** 가 배타 임대해 준다 — 워커끼리, 두 소스끼리, 그리고 + # **여러 프로세스끼리** 같은 IP 를 동시에 쓰거나 태운 IP 를 곧바로 재사용하는 걸 막는다. + # owner 에 PID 를 넣어 프로세스가 늘어도 임대자가 구분된다(죽으면 leased_until 로 자동 회수). + tag = f"p{os.getpid()}-w{i}" + proxy = DecodoProxy(store=store, owner=f"coupang-{tag}") + kr_proxy = DecodoProxy(host=decodo_config.kr_host or None, store=store, owner=f"naver-{tag}") + # 예전엔 seed_offset 으로 워커별 시작 포트를 벌렸지만, 이제 DB 가 LRU(가장 오래 안 쓴 IP)로 + # 배정하므로 프로세스·워커가 몇 개든 알아서 갈린다 — 오프셋 계산이 필요 없다. bot_log = BotDetectionLog() ip_log = IpSessionLog() - suffix = f"_w{i}" if concurrency > 1 else "" + def _pf(source): + """이 워커가 배타적으로 쓸 Chrome 프로필 경로. - def _pf(source): # 워커별 Chrome 프로필 경로(중복 실행 시 ProcessSingleton 충돌 방지) - # [WorkerConfig].profile_dir 를 영속 볼륨으로 두면 재시작해도 cf_clearance 등 쿠키 유지(재웜업 회피). - return f"{worker_config.profile_dir}/lps_{source}{suffix}" + Chrome 은 프로필당 인스턴스 1개만 허용한다(ProcessSingleton). 예전엔 워커 인덱스로만 + 갈랐는데, 그러면 **프로세스가 여러 개일 때 같은 경로를 잡아** 두 번째 프로세스의 검색이 + 전부 죽는다(실측). 파일 락으로 슬롯을 선점해 프로세스가 몇 개든 겹치지 않게 한다. + [WorkerConfig].profile_dir 를 영속 볼륨으로 두면 슬롯 재사용으로 웜 쿠키가 유지된다. + """ + return claim_profile_slot(worker_config.profile_dir, source, worker_index=i) adapters = { "coupang": CoupangAdapter(headless=False, user_data_dir=_pf("coupang"), proxy=proxy, @@ -127,31 +132,27 @@ async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float LOG.w(f"[warmup:{ad.source}] {tries}회 실패(첫 잡에서 재시도): {type(ex).__name__}") -def _proxy_ports_snapshot(adapters) -> tuple[int, int] | None: - """가용 포트 현황 — (최소 가용 수, 전체 포트 수). 프록시 미사용이면 None. +async def _port_pool_status(adapters) -> tuple[dict, tuple[int, int] | None]: + """포트 장부 현황 → (게이트웨이별 상세, (최소 가용, 전체)). - 게이트웨이가 둘(쿠팡=국가무지정 / 네이버=한국)이라 **가장 마른 게이트웨이 기준(min)**으로 본다 - — 한쪽만 고갈돼도 그 소스는 검색을 못 하므로 평균으로 덮으면 안 된다.""" - proxies = {} + 게이트웨이가 둘(쿠팡=국가무지정 / 네이버=한국)이라 **가장 마른 쪽 기준(min)**으로 본다 — + 한쪽만 고갈돼도 그 소스는 검색을 못 하므로 평균으로 덮으면 안 된다. + 가용 수는 DB 장부에서 읽는다(프로세스 로컬 카운터는 다른 프로세스의 차단을 모른다). + """ 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 {} + store = getattr(getattr(ad, "_proxy", None), "_store", None) + if store is None: + continue + try: + detail = await store.snapshot() + except Exception as ex: + LOG.d(f"[ops] 포트 장부 조회 실패(무시): {type(ex).__name__}") + return {}, None + if not detail: + return {}, None + worst = min(detail.values(), key=lambda d: d["available"]) + return detail, (worst["available"], worst["total"]) + return {}, None async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0, adapters=None, alerts=None, ip_log=None): @@ -184,13 +185,12 @@ async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0, adapters 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) + pool, ports = await _port_pool_status(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}} + snap["port_pool"] = pool # {게이트웨이: {held, resting, cooling, available, total}} await alerts.check("proxy_ports_low", avail * 100 <= total * th.ports_low_pct, f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호" + (f" · {pool}" if pool else ""), snap) @@ -291,10 +291,14 @@ 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) + # 포트(=IP 세션) 장부는 DB 다 — 프로세스가 늘어도 모두 같은 장부를 본다. + # 기동 시 게이트웨이별 포트 행을 보장한다(이미 있으면 상태 유지, 덮지 않음). + store = PortLeaseStore() + for gw in {decodo_config.host, decodo_config.kr_host or decodo_config.host}: + if gw: + await store.ensure_ports(gw, 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) + handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history, store=store) all_adapters += worker_adapters bg_tasks.append(asyncio.create_task(_warmup_worker(worker_adapters))) # 챌린지 쿠키 선점(백그라운드) listener = JobListener()