동시 다상품 검색 점검 중 발견. 워커 3개(소유자 6)가 포트 2개/게이트웨이를 두고 경합하는
상황을 실제 코드로 돌리니, 임대를 못 받은 워커가 **남이 쥔 포트를 그대로 집어 같은 IP 로
동시에 요청**했다:
coupang-w1 사용=70002 임대=70002
coupang-w2 사용=70002 임대=None ← 같은 IP 를 둘이 사용
원인은 _port() 의 계산식 폴백이다. 장부 모드에서 acquire 가 None 을 줘도 시간창 계산으로
포트를 하나 골라 돌려줬다. 포트 장부가 존재하는 이유("워커 N개가 같은 IP 에 요청을 몰면 그 IP 가
빨리 탄다" — port_registry.py 도입 배경)를 정면으로 무너뜨리는 경로다. 게다가 하필 **풀이 마른
상태 = IP 가 가장 귀할 때** 발동해, 남은 IP 를 두 배 속도로 태우는 악순환을 만든다.
→ 장부 모드에선 임대한 포트만 쓴다(없으면 None). 못 받으면 AdapterError 로 실패하고 잡이
백오프 후 재시도한다 — 그 사이 쿨다운이 풀린다. 풀 고갈 자체는 proxy_ports_low 가 이미 운다.
→ playwright_proxy() 도 임대가 없으면 예외. 여기서 None 을 돌려주면 **프록시 없이** 브라우저가
떠 서버 공인 IP 로 크롤하게 되는데, 그 IP 가 타면 회전으로 복구할 수 없다.
동시성 점검 결과(포트 20개/게이트웨이, 워커 3개):
정상 24건 동시 성공 24 · 포트 중복 보유 0
풀 고갈 성공 4/6(2건은 정상적으로 실패) · **같은 IP 공유 0**
전면 차단 소각이 어댑터당 2개에서 멈춤(게이트웨이당 6/20) · 브레이커 6/6 트립
테스트 3건 추가(고갈 시 None 반환·남의 포트 미사용 / 임대 없는 playwright_proxy 예외 /
검색이 깔끔히 실패), 전체 256 passed.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
203 lines
10 KiB
Python
203 lines
10 KiB
Python
"""프로세스 간 공유 포트 장부(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)
|
|
|
|
|
|
# ── 종료 시 반납 (2026-08-05) ────────────────────────────────────────────
|
|
async def test_proxy_release_returns_the_port_immediately(store):
|
|
"""워커 종료가 임대를 놓아야 한다 — 안 놓으면 재시작해도 최대 sticky 수명(10분)까지 못 쓴다."""
|
|
from services.search.proxy import DecodoProxy
|
|
|
|
px = DecodoProxy(host=HOST_A, store=store, owner="shutdown-w0")
|
|
px.username, px.password = "u", "p"
|
|
px.port_start, px.port_end = P_START, P_END
|
|
port = await px.ensure_port()
|
|
assert port is not None
|
|
assert (await store.snapshot(HOST_A))[HOST_A]["held"] == 1
|
|
|
|
await px.release()
|
|
assert (await store.snapshot(HOST_A))[HOST_A]["held"] == 0, "종료 후에는 남이 곧바로 쓸 수 있어야 한다"
|
|
assert px._leased is None
|
|
assert await store.acquire(HOST_A, "other-process", 600) is not None
|
|
|
|
|
|
async def test_proxy_release_is_safe_without_a_lease(store):
|
|
"""임대를 못 잡은 채 종료돼도(풀 고갈 등) 예외 없이 지나가야 한다."""
|
|
from services.search.proxy import DecodoProxy
|
|
|
|
px = DecodoProxy(host=HOST_A, store=store, owner="never-leased")
|
|
px.username, px.password = "u", "p"
|
|
px.port_start, px.port_end = P_START, P_END
|
|
await px.release()
|
|
|
|
|
|
# ── 풀 고갈 시 남의 IP 를 빌려 쓰지 않는다 (2026-08-06 회귀) ─────────────
|
|
# 임대 실패 시 계산식으로 포트를 고르는 폴백이 있었다. 그러면 임대를 못 받은 워커가 **남이 쥔
|
|
# 포트를 그대로 집어 같은 IP 에 요청이 겹친다**(실측: 포트 2개·소유자 3인 상황에서 재현).
|
|
# 하필 풀이 마른 상태 = IP 가 가장 귀할 때 벌어져 남은 IP 를 두 배로 태운다.
|
|
|
|
def _proxy_for(store, owner, lo=P_START, hi=P_END):
|
|
from services.search.proxy import DecodoProxy
|
|
px = DecodoProxy(host=HOST_A, store=store, owner=owner)
|
|
px.username, px.password = "u", "p"
|
|
px.port_start, px.port_end = lo, hi
|
|
px.session_minutes = 10
|
|
return px
|
|
|
|
|
|
async def test_exhausted_pool_yields_no_port_instead_of_borrowing(store):
|
|
holders = [_proxy_for(store, f"w{i}") for i in range(P_END - P_START + 1)]
|
|
for px in holders: # 풀을 전부 소진
|
|
assert await px.ensure_port() is not None
|
|
|
|
latecomer = _proxy_for(store, "late")
|
|
assert await latecomer.ensure_port() is None, "빈 포트가 없으면 None 이어야 한다"
|
|
assert latecomer.current_port is None, "계산식 폴백으로 남의 포트를 집으면 안 된다"
|
|
|
|
taken = {px._leased for px in holders}
|
|
assert latecomer._leased not in taken
|
|
|
|
|
|
async def test_no_proxy_config_when_lease_missing(store):
|
|
"""임대가 없는데 playwright_proxy() 가 None 을 주면 **프록시 없이** 브라우저가 떠
|
|
서버 공인 IP 로 크롤한다(그 IP 가 타면 회전으로 복구 불가). 조용히 넘어가면 안 된다."""
|
|
import pytest as _pytest
|
|
px = _proxy_for(store, "no-lease")
|
|
with _pytest.raises(RuntimeError, match="가용 프록시 포트 없음"):
|
|
px.playwright_proxy()
|
|
|
|
|
|
# ── 배타 임대(프로세스 간) ───────────────────────────────────────────────
|
|
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
|