Files
kis_bot/kis_trader/network/market_guard.py
Hwang 61c72a8a4c feat(tests): 신규 키움 웹소켓 조건검색 및 실시간 조건검색 테스트 추가
변경 사항
----
- _test_kiwoom_condition_list.py: 키움 웹소켓 조건검색 '목록조회' 기능을 단독으로 테스트하는 스크립트 추가
- _test_kiwoom_condition_realtime.py: 'momentum' 조건식을 실시간으로 등록하고 초기 매칭 종목 리스트 및 실시간 편입/이탈을 수신하는 테스트 스크립트 추가
- _verify_columnar_bitid.py, _verify_shared_e2e_breakout.py, _verify_shared_e2e.py: 공유 메모리 및 dict 간의 데이터 일관성을 검증하는 테스트 추가

영향
----
- 신규 테스트 스크립트 추가로 키움 웹소켓 API의 기능 검증 및 안정성을 높임
- 기존 기능에 대한 영향 없음

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-06 01:27:00 +09:00

573 lines
22 KiB
Python
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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.
"""
kis_trader/network/market_guard.py — 시장 급락 서킷브레이커
==============================================================
KOSPI/KOSDAQ 종합지수를 주기적으로 폴링해, 다음 조건 충족 시 PANIC 모드 진입:
1) 5분 누적 -N% 하락 (기본 -2%)
2) 일중 누적 -N% 하락 (전일 종가 대비, 기본 -3%)
PANIC 모드 시:
- ``BaseStrategy._scan_and_buy()`` 가 신규 매수를 즉시 차단
- 기존 보유 종목 매도는 정상 동작 (포지션 정리·손실 확대 방지)
해제:
- 5분 누적 +N% 반등 시 자동 해제 (기본 +1%)
- 또는 운영자가 DB ``MARKET_GUARD_ENABLED=false`` 로 수동 해제
영속 (재시작 구멍 방지)
----------------------
- 지수 1분봉: ``ws_candles`` code ``MG0001`` / ``MG1001`` … (timeframe=1, source=mg)
- PANIC 상태: ``kv_store`` ``market_guard.*`` (진입/해제·기동 복원 시 즉시 재판정)
- ``bootstrap_sync()``: 전략 쓰레드 기동 **전** 1회 — DB 봉 복원 → API 1틱 → PANIC 재계산
설계 철학
---------
- 거래소 공식 서킷브레이커 (KOSPI -8%) 보다 훨씬 빨리 반응 → 봇 보호 우선.
- 매수만 차단, 매도는 평소처럼 진행 → 봇이 살아있어야 손절·익절 가능.
- ``MARKET_GUARD_ENABLED=false`` 가 기본값. true 일 때만 지수 API·봉 저장·PANIC 판정.
환경변수 (DB env_config, 매 tick 재조회 → 운영 중 즉시 반영)
------------------------------------------------------------
MARKET_GUARD_ENABLED (true/false, 기본 false)
MARKET_GUARD_5MIN_DROP_PCT (실수, 기본 2.0)
MARKET_GUARD_DAILY_DROP_PCT (실수, 기본 3.0)
MARKET_GUARD_RECOVERY_PCT (실수, 기본 1.0)
MARKET_GUARD_INDEX_CODE ("0001"|"1001"|"both", 기본 "both")
MARKET_GUARD_POLL_SEC (정수, 기본 30)
MARKET_GUARD_INDEX_CANDLE_KEEP_MIN (정수, 기본 20) — 재시작 시 복원할 1분봉 개수
MARKET_GUARD_PERSIST_STATE (true/false, 기본 true) — kv PANIC 영속
"""
from __future__ import annotations
import threading
import time
from collections import deque
from datetime import datetime as dt
from typing import Deque, Dict, List, Optional, Tuple
from ..utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int
from ..utils.logger import get_logger
logger = get_logger("kis_trader.market_guard")
# ws_candles 가상 종목코드: MG + 4자리 지수코드 (예: MG0001=KOSPI)
_WS_CODE_PREFIX = "MG"
_TIMEFRAME_1M = 1
_SOURCE_INDEX = "mg"
KV_PANIC = "market_guard.panic"
KV_PANIC_REASON = "market_guard.panic_reason"
KV_PANIC_SINCE = "market_guard.panic_since_epoch"
KV_PANIC_INDEX = "market_guard.panic_index"
KV_PREV_CLOSE_PREFIX = "market_guard.prev_close."
class MarketGuard:
"""KOSPI/KOSDAQ 지수 폴링 + PANIC 판정 매니저. 독립 백그라운드 쓰레드."""
# 지수 코드 → 사람용 이름
_INDEX_NAMES = {
"0001": "KOSPI",
"1001": "KOSDAQ",
"2001": "KOSPI200",
}
def __init__(self, *, client, db=None):
self.client = client
self.db = db
self._thread: Optional[threading.Thread] = None
self._running = False
self._lock = threading.Lock()
# 지수별 최근 가격 deque [(epoch_sec, value), ...] — 5분 비교용
self._history: Dict[str, Deque[Tuple[float, float]]] = {}
# 지수별 전일 종가 (일중 누적 판정용)
self._prev_close: Dict[str, float] = {}
# PANIC 상태 (전역 — 어떤 지수가 트리거했든 매수 전체 차단)
self._panic = False
self._panic_reason: str = ""
self._panic_since: float = 0.0
self._panic_index: str = ""
self._reload_config()
# ------------------------------------------------------------------
def _db_raw(self):
if self.db is None:
return None
return getattr(self.db, "raw", self.db)
@staticmethod
def _ws_code(index_code: str) -> str:
"""ws_candles 저장용 6자리 코드 (MG0001 …)."""
c = (index_code or "").strip()
return f"{_WS_CODE_PREFIX}{c}"[:20]
@staticmethod
def _candle_minute_key(ts: Optional[float] = None) -> str:
t = dt.fromtimestamp(ts if ts is not None else time.time())
return t.strftime("%Y%m%d%H%M")
@staticmethod
def _minute_end_epoch(candle_time: str) -> Optional[float]:
"""YYYYMMDDHHMM 봉의 종료 시각(epoch, 분 말)."""
ct = (candle_time or "").strip()
if len(ct) < 12:
return None
try:
base = dt.strptime(ct[:12], "%Y%m%d%H%M")
return base.timestamp() + 60.0
except ValueError:
return None
# ------------------------------------------------------------------
def _reload_config(self) -> None:
"""매 tick DB 재조회 (운영 중 임계치 튜닝 즉시 반영)."""
self.enabled = get_env_bool("MARKET_GUARD_ENABLED", False)
self.drop_5min = get_env_float("MARKET_GUARD_5MIN_DROP_PCT", 2.0)
self.drop_daily = get_env_float("MARKET_GUARD_DAILY_DROP_PCT", 3.0)
self.recovery_pct = get_env_float("MARKET_GUARD_RECOVERY_PCT", 1.0)
self.poll_sec = max(5, get_env_int("MARKET_GUARD_POLL_SEC", 30))
self.candle_keep_min = max(5, get_env_int("MARKET_GUARD_INDEX_CANDLE_KEEP_MIN", 20))
self.persist_state = get_env_bool("MARKET_GUARD_PERSIST_STATE", True)
idx_raw = (get_env_from_db("MARKET_GUARD_INDEX_CODE", "both") or "both") \
.strip().lower()
if idx_raw == "both":
self.index_codes = ["0001", "1001"]
else:
codes = [c.strip() for c in idx_raw.split(",") if c.strip()]
self.index_codes = codes if codes else ["0001"]
# ------------------------------------------------------------------
def start(self) -> bool:
"""백그라운드 쓰레드 기동.
``MARKET_GUARD_ENABLED=false`` 면 루프는 대기만(지수 API·봉 저장 없음).
DB 에서 true 로 바꾸면 재시작 없이 다음 주기부터 감시 시작.
"""
if self._thread and self._thread.is_alive():
return True
self._running = True
self._thread = threading.Thread(
target=self._loop, daemon=True, name="MarketGuard"
)
self._thread.start()
logger.info(
"✅ MarketGuard 쓰레드 시작 (enabled=%s, indexes=%s, "
"5min=-%.1f%%, daily=-%.1f%%, recovery=+%.1f%%, poll=%ds, "
"candle_keep=%d분, persist=%s)",
self.enabled, self.index_codes, self.drop_5min, self.drop_daily,
self.recovery_pct, self.poll_sec, self.candle_keep_min,
self.persist_state,
)
return True
def stop(self) -> None:
self._running = False
def bootstrap_sync(self) -> None:
"""기동 시 1회: DB 1분봉·kv PANIC 복원 후 즉시 API 폴링·재판정.
전략 매수 루프보다 **먼저** 호출해야 재시작 직후 매수 구멍을 막을 수 있음.
``MARKET_GUARD_ENABLED=false`` 면 no-op.
"""
self._reload_config()
if not self.enabled:
logger.info(" [MarketGuard] bootstrap 스킵 (ENABLED=false)")
return
self._restore_history_from_db()
self._restore_panic_from_kv()
try:
self._tick()
except Exception as e:
logger.error("[MarketGuard] bootstrap _tick 실패: %s", e)
with self._lock:
panic = self._panic
reason = self._panic_reason
if panic:
logger.warning(
"⛔ [MarketGuard] bootstrap 완료 — PANIC 유지/재진입: %s",
reason,
)
else:
logger.info("📊 [MarketGuard] bootstrap 완료 — 정상 (신규 매수 허용)")
# ------------------------------------------------------------------
# DB 영속
# ------------------------------------------------------------------
def _restore_history_from_db(self) -> None:
"""ws_candles MG* 1분봉 → _history deque 복원 (5분 판정용)."""
raw = self._db_raw()
if raw is None:
return
limit = max(self.candle_keep_min, 12)
with self._lock:
for code in self.index_codes:
ws_code = self._ws_code(code)
try:
rows = raw.get_ws_candles(ws_code, _TIMEFRAME_1M, limit=limit)
except Exception as e:
logger.debug("지수 봉 복원 실패 (%s): %s", ws_code, e)
rows = []
hist: Deque[Tuple[float, float]] = deque(maxlen=120)
for row in rows or []:
ep = self._minute_end_epoch(str(row.get("candle_time", "")))
try:
cl = float(row.get("close", 0) or 0)
except (TypeError, ValueError):
continue
if ep and cl > 0:
hist.append((ep, cl))
if hist:
self._history[code] = hist
pk = f"{KV_PREV_CLOSE_PREFIX}{code}"
try:
pv = raw.get_kv(pk)
if pv:
self._prev_close[code] = float(str(pv).replace(",", ""))
except Exception:
pass
logger.info(
"📂 [MarketGuard] DB 1분봉 복원 — %s",
", ".join(
f"{self._INDEX_NAMES.get(c, c)}={len(self._history.get(c, []))}"
for c in self.index_codes
),
)
def _restore_panic_from_kv(self) -> None:
if not self.persist_state:
return
raw = self._db_raw()
if raw is None:
return
try:
flag = (raw.get_kv(KV_PANIC) or "0").strip()
if flag not in ("1", "true", "True"):
return
reason = (raw.get_kv(KV_PANIC_REASON) or "").strip()
idx = (raw.get_kv(KV_PANIC_INDEX) or "").strip()
since_s = raw.get_kv(KV_PANIC_SINCE)
since = float(since_s) if since_s else time.time()
except Exception as e:
logger.debug("PANIC kv 복원 실패: %s", e)
return
with self._lock:
self._panic = True
self._panic_reason = reason or "kv 복원"
self._panic_index = idx
self._panic_since = since
logger.warning(
"📂 [MarketGuard] kv PANIC 복원 — %s (이후 bootstrap tick 에서 재판정)",
reason or idx,
)
def _persist_panic(self) -> None:
if not self.persist_state:
return
raw = self._db_raw()
if raw is None:
return
with self._lock:
panic = self._panic
reason = self._panic_reason
since = self._panic_since
idx = self._panic_index
try:
raw.set_kv(KV_PANIC, "1" if panic else "0")
raw.set_kv(KV_PANIC_REASON, reason or "")
raw.set_kv(KV_PANIC_SINCE, str(since if panic else 0.0))
raw.set_kv(KV_PANIC_INDEX, idx or "")
except Exception as e:
logger.debug("PANIC kv 저장 실패: %s", e)
def _persist_prev_close(self, code: str, prdy_close: float) -> None:
if not self.persist_state or prdy_close <= 0:
return
raw = self._db_raw()
if raw is None:
return
try:
raw.set_kv(f"{KV_PREV_CLOSE_PREFIX}{code}", str(prdy_close))
except Exception:
pass
def _upsert_index_minute_candle(self, code: str, cur: float, now: float) -> None:
"""폴링 샘플을 1분 OHLC 로 ws_candles 에 적재 (재시작 5분 판정용)."""
raw = self._db_raw()
if raw is None or cur <= 0:
return
ws_code = self._ws_code(code)
minute_key = self._candle_minute_key(now)
try:
latest = raw.get_latest_ws_candle(ws_code, _TIMEFRAME_1M)
except Exception:
latest = None
if latest:
prev_key = str(latest.get("candle_time", "")).strip()
if prev_key and prev_key != minute_key:
try:
o = float(latest.get("open", cur))
h = float(latest.get("high", cur))
l = float(latest.get("low", cur))
c = float(latest.get("close", cur))
raw.upsert_ws_candle(
ws_code, _TIMEFRAME_1M, prev_key,
o, h, l, c, 0, is_confirmed=1, source=_SOURCE_INDEX,
)
except Exception as e:
logger.debug("지수 확정봉 저장 실패: %s", e)
if latest and str(latest.get("candle_time", "")).strip() == minute_key:
try:
o = float(latest.get("open", cur))
h = max(float(latest.get("high", cur)), cur)
l = min(float(latest.get("low", cur)), cur)
raw.upsert_ws_candle(
ws_code, _TIMEFRAME_1M, minute_key,
o, h, l, cur, 0, is_confirmed=0, source=_SOURCE_INDEX,
)
except Exception:
pass
else:
try:
raw.upsert_ws_candle(
ws_code, _TIMEFRAME_1M, minute_key,
cur, cur, cur, cur, 0, is_confirmed=0, source=_SOURCE_INDEX,
)
except Exception:
pass
# ------------------------------------------------------------------
# 외부 API — 전략에서 호출
# ------------------------------------------------------------------
def is_panic(self) -> bool:
"""매수 차단 여부. BaseStrategy._scan_and_buy() 가 매 후보마다 호출."""
with self._lock:
return self.enabled and self._panic
def panic_reason(self) -> str:
"""현재 PANIC 사유 (로그용)."""
with self._lock:
return self._panic_reason
def status(self) -> Dict:
"""heartbeat 로그 등에서 노출용 — 현재 상태 요약."""
with self._lock:
now = time.time()
indexes_status = []
for code in self.index_codes:
name = self._INDEX_NAMES.get(code, code)
hist = self._history.get(code)
prdy = self._prev_close.get(code, 0)
if not hist or prdy <= 0:
indexes_status.append({"name": name, "ready": False})
continue
cur = hist[-1][1]
daily_pct = (cur - prdy) / prdy * 100.0
five = self._lookup_price_at(hist, now - 300)
five_pct = ((cur - five) / five * 100.0) if five else 0.0
indexes_status.append({
"name": name, "ready": True,
"daily_pct": round(daily_pct, 2),
"five_min_pct": round(five_pct, 2),
})
return {
"enabled": self.enabled,
"panic": self._panic,
"reason": self._panic_reason,
"since_min": round((now - self._panic_since) / 60.0, 1) if self._panic else 0.0,
"indexes": indexes_status,
}
# ------------------------------------------------------------------
# 메인 루프
# ------------------------------------------------------------------
def _loop(self) -> None:
last_log_ts = 0.0
was_enabled = False
while self._running:
try:
self._reload_config()
if not self.enabled:
if was_enabled:
logger.info(" [MarketGuard] 비활성 — 지수 폴링·봉 저장 중지")
was_enabled = False
time.sleep(max(10, self.poll_sec))
continue
if not was_enabled:
logger.info("▶ [MarketGuard] 활성화 — 지수 감시·봉 저장 시작")
self._restore_history_from_db()
was_enabled = True
if not self._is_market_hours():
time.sleep(60)
continue
self._tick()
now = time.time()
if now - last_log_ts >= 60:
last_log_ts = now
self._log_status()
except Exception as e:
logger.error("MarketGuard 루프 예외: %s", e)
time.sleep(self.poll_sec)
@staticmethod
def _is_market_hours() -> bool:
"""장 시간 체크 (BaseStrategy.check_market_status 와 동일 정책)."""
if get_env_bool("FORCE_MARKET_OPEN", False):
return True
now = dt.now()
h, m = now.hour, now.minute
return (9 <= h < 15) or (h == 15 and m <= 30)
# ------------------------------------------------------------------
def _tick(self) -> None:
"""모든 감시 지수 1회 폴링 + PANIC 판정."""
now = time.time()
panic_changed = False
for code in self.index_codes:
try:
out = self.client.inquire_index_price(code)
except Exception as e:
logger.debug("지수 조회 실패 (%s): %s", code, e)
continue
if not out:
continue
try:
cur = float(str(out.get("bstp_nmix_prpr", 0)).replace(",", ""))
prdy_close = float(
str(out.get("bstp_nmix_prdy_clpr", 0)).replace(",", "")
)
except (ValueError, AttributeError):
continue
if cur <= 0 or prdy_close <= 0:
continue
with self._lock:
self._prev_close[code] = prdy_close
hist = self._history.setdefault(code, deque(maxlen=120))
hist.append((now, cur))
cutoff = now - 600
while hist and hist[0][0] < cutoff:
hist.popleft()
prev_panic = self._panic
self._evaluate_panic(code, cur, prdy_close, now)
if self._panic != prev_panic:
panic_changed = True
self._persist_prev_close(code, prdy_close)
self._upsert_index_minute_candle(code, cur, now)
if panic_changed:
self._persist_panic()
def _evaluate_panic(
self, code: str, cur: float, prdy_close: float, now: float,
) -> None:
"""단일 지수 기준 PANIC 진입/해제 판정. _lock 보유 상태 가정."""
name = self._INDEX_NAMES.get(code, code)
daily_chg = (cur - prdy_close) / prdy_close * 100.0
hist = self._history.get(code, deque())
five_min_ago_price = self._lookup_price_at(hist, now - 300)
if five_min_ago_price:
five_min_chg = (cur - five_min_ago_price) / five_min_ago_price * 100.0
else:
five_min_chg = 0.0
if not self._panic:
reasons: List[str] = []
if daily_chg <= -abs(self.drop_daily):
reasons.append(f"일중 {daily_chg:+.2f}% (한계 -{self.drop_daily:.1f}%)")
if five_min_chg <= -abs(self.drop_5min):
reasons.append(f"5분 {five_min_chg:+.2f}% (한계 -{self.drop_5min:.1f}%)")
if reasons:
self._panic = True
self._panic_index = code
self._panic_reason = f"{name} " + " / ".join(reasons)
self._panic_since = now
logger.warning(
"🚨 [MarketGuard] PANIC 진입: %s — 모든 전략 신규 매수 차단",
self._panic_reason,
)
else:
if code != self._panic_index:
return
if five_min_chg >= abs(self.recovery_pct):
duration_min = (now - self._panic_since) / 60.0
logger.warning(
"✅ [MarketGuard] PANIC 해제: %s 5분 %+.2f%% 반등 "
"(지속 %.1f분, 매수 재개)",
name, five_min_chg, duration_min,
)
self._panic = False
self._panic_reason = ""
self._panic_index = ""
self._panic_since = 0.0
@staticmethod
def _lookup_price_at(
hist: Deque[Tuple[float, float]], target_ts: float,
) -> Optional[float]:
"""deque 에서 target_ts 이전(≤) 시각 중 가장 최근 가격 반환."""
best = None
for ts, val in hist:
if ts <= target_ts:
best = val
else:
break
return best
def _log_status(self) -> None:
"""60초에 한 번씩 평상시/PANIC 상태 로그."""
with self._lock:
if self._panic:
duration_min = (time.time() - self._panic_since) / 60.0
logger.warning(
"⛔ [MarketGuard] PANIC 모드 진행 중 — %s (지속 %.1f분)",
self._panic_reason, duration_min,
)
return
parts = []
now = time.time()
for code in self.index_codes:
name = self._INDEX_NAMES.get(code, code)
hist = self._history.get(code)
prdy = self._prev_close.get(code, 0)
if not hist or prdy <= 0:
parts.append(f"{name}=대기")
continue
cur = hist[-1][1]
daily = (cur - prdy) / prdy * 100.0
five = self._lookup_price_at(hist, now - 300)
five_chg = ((cur - five) / five * 100.0) if five else 0.0
parts.append(f"{name} 일중{daily:+.2f}% 5분{five_chg:+.2f}%")
if parts:
logger.info("📊 [MarketGuard] 정상 — %s", " | ".join(parts))