o2o-negosium-original/lps/crud/port_lease.py
민헌 3b345929d4 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.
2026-08-05 11:44:54 +09:00

147 lines
7.1 KiB
Python

"""프록시 포트(=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)