o2o-negosium-original/lps/worker_main.py
민헌 377389f495 fix(lps): IP 로테이션 안정성 검수 — 예산이 안 먹던 근본 원인 + 차단 시 풀 소각 차단
크롤·IP 로테이션을 검수하며 찾은 결함을 순서대로 고쳤다. 의심 지점은 모두 실제 코드 경로로
재현해 확인했다(브라우저·네트워크만 mock, 프록시·DB 장부는 실물).

**① 요청 예산이 사실상 발화하지 않았다 (핵심)**
유휴 정리(close_if_idle, 120s)는 브라우저만 닫고 임대는 두는데, 재기동 때마다 _ip_requests 를
0 으로 되돌렸다. 게다가 ensure_port 의 renew 가 임대 만료를 계속 뒤로 민다 — 검색이 유휴
임계보다 뜸하고 임대(10분)보다 잦으면 **한 IP 에 영원히 고정**된다(실측: 6회 검색이 전부 같은
포트·ip_req#1). 연속 검색에서는 정상 동작해 부하 테스트로는 안 잡히고, 수동 트리거처럼
드문드문한 실사용 패턴에서만 깨진다.
파급이 하나 더 있다 — bot_detection.ip_request_no 가 항상 1 로 찍혀, operations.md 가 명시한
'1 위주면 IP 평판 / 2 이상이면 예산 하향' 진단이 통째로 무너진다. 과거 "전량 ip_req#1 이라
IP 평판 문제" 결론은 이 착시일 수 있다(문서에 경고 추가).
→ IP 세션 상태를 브라우저 수명과 분리. **포트가 실제로 바뀔 때만** 리셋한다(_begin_ip_session).
   세션 종료 기록도 포트 변경·최종 close 시점으로 옮겼다(idle 사유 소멸).

**② 환경 차단이면 회복 못 하는데 풀을 계속 태웠다**
쿠팡은 fatal 마커가 없어 컨테이너 차단 같은 '회전 무효' 상황을 구분 못 했다. 실측으로
웜업 6포트 + 잡 1건당 6포트를 30분 쿨다운에 묶어 **잡 16건이면 100포트 고갈**. 실제 장부에도
9분간 11포트 연속 소각 이력이 남아 있다(gate 사용 21 / 소각 14).
→ 서킷브레이커: **서로 다른 IP 가 연속 3개 모두 첫 요청부터** 막히면 IP 문제가 아니라고 판정,
   태우기를 멈추고 fatal 로 알린다(env_block 마커 → 기존 fatal_block 알림이 집계).
   같은 IP 반복 차단·뒤쪽 요청 차단은 세지 않는다. 성공 1회로 자동 해제(타이머 불필요).
   결과: 전면 차단 시 소각이 판정 근거 2개에서 멈춘다(웜업 6→0, 잡 6→0).

**③ 종료가 임대를 반납하지 않았다**
close() 후에도 leased_until(최대 10분)까지 그 IP 를 아무도 못 썼다 — 재시작이 잦을수록 가용
풀이 줄었다. DecodoProxy.release() 추가, close() 에서만 호출(유휴 정리는 웜 쿠키·예산 유지를
위해 그대로 둔다).

**④ 시간창 재기동이 IP 를 안 바꿨다** — 로그만 'IP 회전'이었고 renew 로 같은 포트를 붙잡았다.
sticky 수명이 끝나면 같은 포트라도 IP 가 바뀌므로 명시적으로 놓아준다.

**⑤ '검색결과 없음'을 차단으로 오인해 IP 를 태울 수 있었다**
네이버 무결과 페이지 크기는 실측된 적이 없는데 short_html 폴백이 이를 차단으로 본다.
확신도로 대응을 갈랐다 — 알려진 마커만 태우고/서킷브레이커에 세고, 미지의 짧은 HTML 은
회전·재시도까지만. 판단 근거는 bot_detection 에 계속 쌓이므로 나중에 임계를 실측할 수 있다.

**⑥** available_ports() 가 장부 모드에서 늘 최대값을 반환하는 점을 문서화(관측 경로는 미사용).
세션 마감을 멱등하게 만들어 close() 중복 호출 시 이중 기록 방지.

테스트 14건 추가(전체 251 passed). mock 하니스도 실물을 타도록 고쳤다 — 회전 시 포트가 실제로
바뀌고, 재기동 판단·IP 세션 경계는 실제 코드를 그대로 쓴다(고정 포트 mock 은 이 버그를 못 봤다).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-05 17:09:01 +09:00

395 lines
25 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# LPS 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일).
# python worker_main.py
# WORKER_CONCURRENCY=3 python worker_main.py # 상품 3개 동시 검색(권장 2~3, 로컬)
#
# 브라우저 어댑터는 컨텍스트당 직렬(lock)이라, 진짜 병렬을 위해 **워커마다 자기 브라우저 세트**를 준다:
# 프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP). 동시성 N → 최대 4×N Chrome.
import asyncio
import os
import signal
import time
from common.alerts import AlertManager
from common.database.db_session_manager import DB_SESSION_MNG
from common.logger import LOG
from config.server_configs import web_server_config, openai_config, decodo_config, worker_config, alert_config
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.browser_base import ENV_BLOCK_MARKER
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
from services.search.esm.adapter import EsmAdapter
from services.search.st11.adapter import ElevenStAdapter
from services.ai.similarity import SimilarityJudge
from services.ai.keyword import KeywordGenerator
from worker.handlers import build_search_handler
from worker.notify import JobListener
from worker.runner import Worker, run_reaper
LOG.SetPrefix(f"{web_server_config.server_name}-worker")
# 오픈마켓 폴백(G마켓·옥션·11번가)은 **기본 비활성** — 2026-07-10 협의 결정.
# 실측상 크롤 몰이 최종 최저가를 바꾼 적이 없고(0회), 검색당 최대 15s + 프록시 대역폭의
# 대부분을 차지해 로직에서 제외했다(코드·테스트는 유지, 핸들러는 빈 폴백을 정상 처리).
# 재가동: [WorkerConfig].fallbacks = ["gmarket", "auction", "st11"] (일부만도 가능) —
# 켜기 전 라이브 스모크로 셀렉터 드리프트 점검. 배경은 docs/decision-openmarket-crawler.md.
_FALLBACK_SOURCES = ("gmarket", "auction", "st11")
def _enabled_fallbacks() -> list[str]:
names = [s.strip() for s in worker_config.fallbacks if s.strip()]
unknown = [n for n in names if n not in _FALLBACK_SOURCES]
if unknown:
LOG.w(f"[WorkerConfig].fallbacks 무시된 값: {unknown} (가능: {list(_FALLBACK_SOURCES)})")
return [n for n in names if n in _FALLBACK_SOURCES]
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포트를 워커 수로 균등 분할해 시작점을 벌린다.
# 두 프록시는 **게이트웨이가 다르다** — 쿠팡=국가 무지정, 네이버=한국 전용.
# 해외 residential IP 로는 네이버 msearch 가 즉시 하드차단된다(2026-08-05 실측:
# 같은 포트 10091 에서 gate=차단 2,641B / kr=정상 14건). kr_host 가 비면 gate 로
# 떨어지고 그때는 네이버가 막힌다(로그의 [naver][BOT-DETECTED] 로 드러남).
# 포트는 **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()
def _pf(source):
"""이 워커가 배타적으로 쓸 Chrome 프로필 경로.
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,
on_detect=bot_log.record, on_session_end=ip_log.record),
# 네이버는 오픈API(shop.json)가 2026-07-31 종료돼 404 SE05 만 돌아온다 → 모바일 크롤로 전환.
# 쿠팡과 같은 프록시를 물린다: 단일 IP 로 연달아 두드리면 WTM 캡차로 넘어간다(실측 —
# 스파이크에서 십수 회 만에 홈 IP 가 탔다). 회전·예산·차단기록은 베이스가 처리한다.
"naver": NaverShopAdapter(headless=False, user_data_dir=_pf("naver"), proxy=kr_proxy,
ip_request_budget=decodo_config.naver_ip_request_budget,
on_detect=bot_log.record, on_session_end=ip_log.record),
}
# 폴백은 기본 비활성(LPS_FALLBACKS 로 켬 — 상단 주석 참고). 켤 땐 데드라인이 상한이라
# 봇감지 재시도(챌린지 대기 2배)를 끈다(max_block_retries=0) — 빠르게 포기·스킵.
fallback_adapters = {}
for name in _enabled_fallbacks():
if name == "st11":
fallback_adapters[name] = ElevenStAdapter(headless=False, user_data_dir=_pf("st11"), proxy=proxy,
on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0)
else:
fallback_adapters[name] = EsmAdapter(name, headless=False, user_data_dir=_pf(name), proxy=proxy,
on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0)
# AI 도 워커별 인스턴스 — 공유 상태(last_usage) 경합 원천 제거
judge = SimilarityJudge() if has_openai else None
keyword_gen = KeywordGenerator() if has_openai else None
handler = build_search_handler(
adapters, judge=judge, keyword_gen=keyword_gen,
neg_cache=neg_cache, history=history,
fallback_adapters=fallback_adapters,
ai_model=openai_config.model,
proxy_cost_per_gb=decodo_config.cost_per_gb,
)
return handler, list(adapters.values()) + list(fallback_adapters.values())
# 웜업 대상 = 브라우저로 긁는 소스 전부. 챌린지 쿠키 선점이 목적이지만, 동시에
# **기동 직후 크롤 가능 여부를 확인하는 프리플라이트**이기도 하다 — 여기서 실패하면
# 그 소스는 이 환경에서 아예 못 긁는다는 뜻이라, 잡이 쌓여 DEAD 될 때까지 기다리지 않고 바로 알린다.
# (실측 배경: 컨테이너 워커는 쿠팡·네이버가 모두 차단되는데 프로세스는 healthy 라 조용히 0건이 된다)
_WARMUP_SOURCES = ("gmarket", "auction", "coupang", "naver")
async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float = 60.0, alerts=None):
"""워커의 크롤 소스를 미리 한 번 긁어 (1) 챌린지 쿠키 확보 (2) 크롤 가능 여부 확인.
콜드 비용을 시작 시 몰아, 이후 실 작업은 웜(빠름). 백그라운드로 돌려 잡 처리를 막지 않는다.
나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다.
시도당 타임아웃 필수 — 웜업은 search 중 어댑터 락을 쥐므로, 여기서 행하면 그 워커의
모든 실 검색이 락 대기로 함께 멈춘다(2026-07-10 부하테스트에서 15분 행 실측)."""
for ad in worker_adapters:
if ad.source not in _WARMUP_SOURCES:
continue
for attempt in range(tries):
try:
await asyncio.wait_for(ad.search("생수", limit=1), timeout=attempt_timeout)
LOG.i(f"[warmup:{ad.source}] 통과·쿠키 확보 (시도 {attempt + 1})")
if alerts is not None:
await alerts.check(f"warmup_{ad.source}", False, f"{ad.source} 크롤 정상")
break
except Exception as ex:
if attempt < tries - 1:
ad._rotate_ip(f"웜업 재시도({type(ex).__name__}) — 새 IP")
else:
msg = (f"[warmup:{ad.source}] {tries}회 모두 실패({type(ex).__name__}: {str(ex)[:80]}) — "
f"이 환경에서 {ad.source} 크롤이 막혔을 수 있습니다. "
f"호스트에서는 되는데 컨테이너에서만 막히는 사례가 있으니 실행 환경을 확인하세요")
LOG.w(msg)
if alerts is not None:
await alerts.check(f"warmup_{ad.source}", True, msg)
def _fatal_markers(adapters) -> list[str]:
"""'회전 무효' 마커 목록 — 알림이 이 마커만 세도록 모은다.
두 종류를 함께 센다: 어댑터가 HTML 에서 알아보는 구조적 차단 마커와, 연속 패턴으로
판정하는 환경 차단(서킷브레이커 트립, ENV_BLOCK_MARKER). 둘 다 IP 회전으로는 회복되지
않아 사람이 환경/설정을 고쳐야 한다는 점에서 같은 계열이다.
"""
out = [ENV_BLOCK_MARKER]
for ad in (adapters or []):
out += list(getattr(ad, "fatal_block_markers", ()) or ())
return sorted(set(out))
async def _port_pool_status(adapters) -> tuple[dict, tuple[int, int] | None]:
"""포트 장부 현황 → (게이트웨이별 상세, (최소 가용, 전체)).
게이트웨이가 둘(쿠팡=국가무지정 / 네이버=한국)이라 **가장 마른 쪽 기준(min)**으로 본다 —
한쪽만 고갈돼도 그 소스는 검색을 못 하므로 평균으로 덮으면 안 된다.
가용 수는 DB 장부에서 읽는다(프로세스 로컬 카운터는 다른 프로세스의 차단을 모른다).
"""
for ad in (adapters or []):
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):
"""워커 헬스 하트비트 + 임계 알림. 주기적으로 (1) 하트비트 파일 갱신(Docker HEALTHCHECK 가
행/좀비 워커 감지) (2) 큐/차단/DB풀/소스별 실패 지표 점검 → AlertManager 로 발화
(룰별 쿨다운으로 스팸 방지, 조건 해소 시 회복 알림)."""
hb_path = worker_config.heartbeat_file
th = alert_config # 임계값은 [AlertConfig] 섹션이 소스(docs/operations.md 표)
alerts = alerts or AlertManager(origin="worker")
ip_log = ip_log or IpSessionLog()
while not stop.is_set():
try:
with open(hb_path, "w") as f:
f.write(str(int(time.time()))) # 하트비트(mtime) — HEALTHCHECK 가 신선도 확인
except Exception:
pass
try:
snap = await queue.ops()
snap["blocks_1h"] = await bot_log.recent_count(60)
pool = DB_SESSION_MNG.pool_status()
snap["pool_pct"] = pool["pct"]
await alerts.check("dead", snap["dead_1h"] >= th.dead_1h, f"DEAD 1h={snap['dead_1h']}", snap)
await alerts.check("blocks", snap["blocks_1h"] >= th.blocks_1h, f"차단 1h={snap['blocks_1h']}", snap)
# 구조적 차단(회전 무효) — 1건만 나와도 알린다. 방치하면 그 소스는 계속 0건이다.
fatal = _fatal_markers(adapters)
if fatal:
snap["fatal_blocks_1h"] = await bot_log.recent_count_by_marker(fatal, 60)
await alerts.check("fatal_block", snap["fatal_blocks_1h"] > 0,
f"구조적 차단 1h={snap['fatal_blocks_1h']} (마커 {fatal}) — "
f"IP 회전으로 회복 불가. 게이트웨이 국가 설정([DecodoConfig].kr_host) 확인", snap)
await alerts.check("queue_lag", snap["oldest_pending_sec"] >= th.queue_lag_sec, f"큐지연={snap['oldest_pending_sec']}s", snap)
await alerts.check("stuck", snap["stuck_running"] > 0, f"stuck={snap['stuck_running']}", snap)
await alerts.check("db_pool", pool["pct"] >= th.pool_pct,
f"DB 풀 포화 {pool['pct']}% (checked_out {pool['checked_out']}/{pool['capacity']})", snap)
await alerts.check("deadline", snap["deadline_1h"] >= th.deadline_1h,
f"잡 데드라인 강제종료 1h={snap['deadline_1h']} — 크롤 행 반복 신호", snap)
await alerts.check("cost", snap["cost_1h_usd"] >= th.cost_1h_usd,
f"검색원가 1h=${snap['cost_1h_usd']} — 비용 폭주(리소스차단 풀림·재시도 루프) 점검", snap)
# 가용 프록시 포트 고갈 — 쿨다운 격리 누적. blocks_1h 보다 먼저 우는 대규모 차단 조기 신호.
pool, ports = await _port_pool_status(adapters)
if ports:
avail, total = ports
snap["proxy_ports_avail"], snap["proxy_ports_total"] = avail, total
if pool:
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)
# 예산 누수 — 요청 예산을 지켰는데도 차단된 IP 세션 발생.
# 처방은 '몇 번째 요청에서 막혔나'로 갈린다:
# ip_req#1 위주 → 새 IP 첫 요청부터 차단 = IP 평판 문제. 예산을 낮춰도 소용없다.
# ip_req#2~ 위주 → 같은 IP 로 너무 많이 긁은 것 = 예산 하향이 유효.
# 예산 하향만 권하면 오진을 부른다(2026-07-28 배포서버 조사에서 전량 ip_req#1 이었다).
block_sessions = (await ip_log.recent_stats(360)).get("block", 0)
snap["block_sessions_6h"] = block_sessions
await alerts.check(
"budget_leak", block_sessions >= th.block_sessions_6h,
f"예산 회전에도 차단된 IP 세션 6h={block_sessions}"
f"bot_detection.ip_request_no 분포 확인(1 위주면 IP 평판/프록시 대역, 2 이상이면 ip_request_budget 하향)",
snap,
)
# 소스별 장기 실패 — 최근 30분간 시도는 있는데 성공이 0건(쿼터 소진·셀렉터 드리프트·전면 차단 신호)
per_source: dict[str, list[int]] = {}
for ad in (adapters or []):
tries, ok = ad.recent_stats(1800)
agg = per_source.setdefault(ad.source, [0, 0])
agg[0] += tries
agg[1] += ok
for src, (tries, ok) in per_source.items():
await alerts.check(f"source_fail:{src}", tries >= th.source_fail_30m and ok == 0,
f"{src} 최근 30분 {tries}회 시도·성공 0건", snap)
except Exception as ex:
LOG.e_no_callstack(f"[ops-monitor] {type(ex).__name__}: {ex}")
try:
await asyncio.wait_for(stop.wait(), timeout=interval)
except asyncio.TimeoutError:
pass
async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: float = 30.0,
close_timeout: float = 60.0):
"""유휴 브라우저 정리 루프 — 일정 시간 검색 없는 어댑터의 Chrome 을 닫아 메모리를 회수한다.
쿠키는 user_data_dir 에 남아, 다음 검색 때 재기동해도 (같은 IP면) 웜 유지.
순차 순회라 close 1건에도 타임아웃을 건다 — 한 어댑터의 close 행이 루프 전체를 멈춰
다른 워커의 브라우저까지 못 닫게 되는 것을 실측(2026-07-10 부하테스트)했다."""
while not stop.is_set():
try:
await asyncio.wait_for(stop.wait(), timeout=interval)
except asyncio.TimeoutError:
pass
for ad in adapters:
close_if_idle = getattr(ad, "close_if_idle", None)
if close_if_idle is None: # 네이버(httpx) 등 브라우저 없는 어댑터는 정리 대상 아님
continue
try:
await asyncio.wait_for(close_if_idle(idle_sec), timeout=close_timeout)
except asyncio.TimeoutError:
LOG.w(f"[browser-reaper] {getattr(ad, 'source', '?')} 정리 {close_timeout:.0f}s 초과 — 취소·스킵(다음 틱 재시도)")
except Exception as ex:
LOG.e_no_callstack(f"[browser-reaper] 정리 실패(무시): {ex}")
async def main(concurrency: int = 1):
queue = JobQueue()
neg_cache, history = NegativeCache(), PriceHistory() # DB 기반 — 워커 공유 안전
has_openai = bool(openai_config.api_key)
# 시작 프리플라이트: DECODO 게이트가 살아있는지(인증) 대표 프록시로 1회 확인. 포트는 워커별로 각자 잡음.
probe = DecodoProxy()
LOG.i(f"DECODO 프록시: {'ON(sticky ' + str(probe.session_minutes) + '분 회전)' if probe.enabled else 'OFF(미설정)'}")
if probe.enabled:
egress_ip, egress_port = await probe.healthcheck()
LOG.i(f"DECODO 프리플라이트 OK — egress IP {egress_ip} (port {egress_port})") if egress_ip \
else LOG.w("DECODO 프리플라이트 실패 — 살아있는 포트를 못 찾음(런타임 회전으로 재시도)")
fb = _enabled_fallbacks()
LOG.i(f"AI(판정+검색어생성): {'ON' if has_openai else 'OFF(키 없음)'} · "
f"오픈마켓 폴백: {', '.join(fb) if fb else 'OFF(기본 — LPS_FALLBACKS 로 활성화)'}")
stop = asyncio.Event()
listeners: list[JobListener] = []
tasks: list[asyncio.Task] = []
bg_tasks: list[asyncio.Task] = [] # 웜업 등 백그라운드(짧게 끝남, gather 대상 아님)
all_adapters = []
# ── graceful shutdown: SIGINT(Ctrl+C)/SIGTERM(docker stop) → stop 이벤트 ──
# asyncio.run 기본 동작(SIGINT=메인 태스크 즉시 cancel)은 하던 잡을 도중에 끊어
# RUNNING 인 채 lease 만료(120s)까지 묶어둔다. 대신 stop 을 set 해 "새 잡은 안 받고,
# 하던 잡은 마무리"로 종료한다. 같은 신호를 한 번 더 받으면 강제 종료(태스크 취소).
def _request_stop(sig_name: str):
if not stop.is_set():
LOG.i(f"{sig_name} 수신 — graceful 종료: 새 잡 중단, 하던 잡 마무리 (한 번 더 = 강제 종료)")
stop.set()
for t in bg_tasks: # 웜업은 선택 작업 — 즉시 취소해 어댑터 락을 비운다
t.cancel()
else:
LOG.w(f"{sig_name} 재수신 — 강제 종료(실행 중 잡은 lease 만료 후 reaper 가 재큐)")
for t in tasks:
t.cancel()
loop = asyncio.get_running_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
loop.add_signal_handler(sig, _request_stop, sig.name)
# 잡 1건 데드라인 — 정상 검색은 폴백 포함 수분 내 끝난다(실측 15~22s). 크롤 행 실측(15분) 대비 상한.
job_deadline = worker_config.job_deadline_sec
# 알림은 웜업과 ops 모니터가 공유한다 — 같은 룰이 양쪽에서 중복 발화하지 않도록.
alerts = AlertManager(origin="worker")
# 포트(=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, store=store)
all_adapters += worker_adapters
bg_tasks.append(asyncio.create_task(_warmup_worker(worker_adapters, alerts=alerts))) # 쿠키 선점 + 크롤 프리플라이트
listener = JobListener()
await listener.start()
listeners.append(listener)
worker = Worker(f"worker-{i}", queue, handler, job_deadline_sec=job_deadline)
tasks.append(asyncio.create_task(worker.run(listener, stop)))
tasks.append(asyncio.create_task(run_reaper(queue, stop)))
tasks.append(asyncio.create_task(run_browser_reaper(all_adapters, stop))) # 유휴 브라우저 정리
tasks.append(asyncio.create_task(run_ops_monitor(queue, BotDetectionLog(), stop, adapters=all_adapters, alerts=alerts))) # 하트비트 + 임계 알림
LOG.i(f"LPS 워커 {concurrency}개 + reaper + 브라우저정리 + ops모니터(하트비트/알림) 기동 (워커별 세트 · 상품 {concurrency}개 동시)")
# 종료 유예: stop 후 하던 잡이 이 시간 안에 끝나면 자연 종료, 초과하면 강제 취소.
# docker stop 을 쓰면 compose 의 stop_grace_period 를 이보다 길게 잡아야 SIGKILL 전에 마무리된다.
grace = worker_config.shutdown_grace_sec
gathered = asyncio.gather(*tasks)
stop_waiter = asyncio.create_task(stop.wait())
try:
await asyncio.wait({gathered, stop_waiter}, return_when=asyncio.FIRST_COMPLETED)
if gathered.done():
gathered.result() # 워커/리퍼가 예외로 죽은 경우 → 전파(finally 가 정리 후 종료)
else:
# 종료 신호 경로 — 워커 루프들이 stop 을 보고 하던 잡을 마친 뒤 스스로 끝나길 기다린다
try:
await asyncio.wait_for(gathered, timeout=grace)
LOG.i("graceful 종료 — 모든 워커가 하던 잡을 마무리함")
except asyncio.TimeoutError:
LOG.w(f"종료 유예 {grace:.0f}s 초과 — 남은 태스크 강제 취소(잡은 lease 만료 후 재큐)")
except asyncio.CancelledError: # 신호 재수신(강제 종료)로 태스크가 취소된 경우
LOG.w("강제 종료 — 남은 리소스 정리 후 종료")
finally:
stop.set()
stop_waiter.cancel()
for t in (*tasks, *bg_tasks):
t.cancel()
# 취소 완주를 기다린 뒤 정리 — 실행 중 태스크가 브라우저/커넥션을 쓰는 채로 닫지 않게
await asyncio.gather(gathered, stop_waiter, *bg_tasks, return_exceptions=True)
for listener in listeners:
try:
await listener.close()
except Exception as ex:
LOG.e_no_callstack(f"[shutdown] 리스너 정리 실패(무시): {ex}")
for adapter in all_adapters: # 항목별 격리 — 하나가 실패해도 나머지 Chrome 은 닫는다
try:
await adapter.close()
except Exception as ex:
LOG.e_no_callstack(f"[shutdown] {getattr(adapter, 'source', '?')} 정리 실패(무시): {ex}")
LOG.i("LPS 워커 종료 완료")
if __name__ == "__main__":
# 동시성은 [WorkerConfig].concurrency 가 소스 — WORKER_CONCURRENCY env 는 실행 스크립트의
# 대화형 입력 전용 임시 override(설정 관리는 toml 하나로, 2026-07-13 협의).
_conc = int(os.environ.get("WORKER_CONCURRENCY", "0")) or worker_config.concurrency
asyncio.run(main(_conc))