feat(lps): 포트(IP) 장부를 DB 로 올려 멀티 프로세스에서 안전하게 — 프로필 슬롯도 분리

한 DECODO 계정을 여러 워커 프로세스가 나눠 쓰는 전제로 전환한다. 인메모리 장부는
프로세스마다 따로라 (1) 같은 IP 를 동시에 잡고 (2) 한쪽이 태운 IP 를 다른 쪽이 곧바로
집으며 (3) 재시작하면 쿨다운이 통째로 사라졌다.

proxy_port 테이블 = 단일 진실. 상태는 세 시각으로만 표현한다(leased/rest/cooldown_until).
- acquire: 한 UPDATE 안에서 FOR UPDATE SKIP LOCKED 로 후보를 잠그고 임대까지 끝낸다
  (잡 큐와 같은 방식 — SELECT 후 UPDATE 로 나누면 그 틈에 다른 프로세스가 같은 행을 집는다)
- 회전은 LRU(last_used_at). 프로세스가 몇 개든 '가장 오래 안 쓴 IP'를 집으므로 전체가
  자연히 한 바퀴씩 돈다 → 프로세스별 seed_offset 계산 제거
- 죽은 프로세스 회수: leased_until 만료로 자동 복귀(별도 reaper 불필요)
- 차단·휴식은 전역이라 재시작해도 유지된다

DB 왕복은 비동기라 검색 루프(동기)에서 곧바로 못 한다 → 회전·차단을 pending 에 적어두고
ensure_port(브라우저 재기동 직전, async)에서 한 번에 flush. _close_ctx 에서도 flush 해
종료 시 유실(=태운 IP 를 남이 그대로 집는 상황)을 막는다.

**프로필 슬롯**(services/search/profile_slot): Chrome 은 user_data_dir 당 1 인스턴스다.
예전엔 워커 인덱스로만 갈라서 프로세스 2개면 같은 경로를 잡아 두 번째가 통째로 죽었다
(실측: 잡 3건 중 2건 DEAD, TargetClosedError). 파일 락으로 슬롯을 선점한다 — PID 경로가
아니라 슬롯이라 재시작 시 재사용돼 웜 쿠키(cf_clearance·Akamai)를 버리지 않는다.

검증: 프로세스 2개 동시 acquire 20회 → 중복 배정 0건. 워커 2프로세스 e2e → 잡 3건 모두
DONE(네이버가 삼다수 최저가 획득 8,960 < 13,200). 테스트 14건 추가, 전체 217 passed.
This commit is contained in:
민헌 2026-08-05 11:44:54 +09:00
parent 495e16b951
commit 3b345929d4
10 changed files with 582 additions and 66 deletions

View File

@ -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"),
)

146
lps/crud/port_lease.py Normal file
View File

@ -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)

View File

@ -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);

View File

@ -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()

View File

@ -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

View File

@ -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):

View File

@ -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)

View File

@ -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

View File

@ -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] # 마지막 슬롯 공유(검색 중단보다 낫다)

View File

@ -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()