""" kiwoom_ws.py — 키움 WebSocket 실시간 시세 캐시 (시세 마이그레이션 검증용) ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ 목적 ---- KIS WS(41 한도)의 시세를 키움 WS(100 한도)로 옮기기 전에, 둘을 동시에 돌려 가격 일치성을 검증하기 위한 키움 WebSocket 클라이언트. 설계 원칙 (KIS WS 와 동일) -------------------------- - ``Dict[str, Dict]`` 메모리 캐시 (락만 잠그고 마이크로초 read/write) - ``get_price(code)`` 인터페이스를 KIS WS 와 100% 동일 포맷으로 제공 → 봇 코드 재사용성 100% - ``KiwoomTokenManager`` 싱글톤 (kis_ws.py 안에 있음) 재사용 - 재연결 백오프, approval 갱신, 종목 등록/해지 키움 WebSocket 스펙 ------------------- URL 실전: wss://api.kiwoom.com:10000/api/dostk/websocket URL 모의: wss://mockapi.kiwoom.com:10000/api/dostk/websocket [프로토콜] 1. 연결 후 LOGIN: {"trnm": "LOGIN", "token": ""} 2. 등록: {"trnm": "REG", "grp_no": "1", "refresh": "1", "data": [{"item": ["005930"], "type": ["0B"]}]} 3. 해지: {"trnm": "REMOVE","grp_no": "1", "data": [{"item": ["005930"], "type": ["0B"]}]} 4. PING: {"trnm": "PING"} → 30초마다 서버가 보냄, 그대로 echo [주식체결 0B 메시지 FID] 10 = 현재가(체결가; 부호 → 음수=하락) 11 = 전일대비 12 = 등락률 13 = 누적거래량 14 = 누적거래대금 15 = 거래량(체결량) 16 = 시가 17 = 고가 18 = 저가 20 = 체결시간(HHMMSS) (키움 OpenAPI+ FID 정의 — 종목별 동일) 설치 ---- ``pip install websocket-client`` (KIS WS 와 공용) 토글 (DB env_config) -------------------- ``WS_PROVIDER`` (기본 ``kis_only`` — 키움 WS 미기동) ``KIWOOM_WS_URL_REAL`` / ``KIWOOM_WS_URL_MOCK`` (URL 재정의용) """ from __future__ import annotations import json import logging import threading import time from typing import Dict, Optional, Set logger = logging.getLogger("KiwoomWebSocket") # ── 환경변수 헬퍼 (KIS WS 와 동일 패턴, fallback 포함) ──────────────────── try: from kis_long_ver1 import get_env_from_db, get_env_int # noqa: F401 except ImportError: try: from kis_short_ver2 import get_env_from_db, get_env_int # noqa: F401 except ImportError: def get_env_from_db(key, default=""): # type: ignore[misc] return default def get_env_int(key, default): # type: ignore[misc] return default # ────────────────────────────────────────────────────────────────────── # 메인 클래스 # ────────────────────────────────────────────────────────────────────── class KiwoomWebSocketPriceCache: """키움 실시간 체결가 (0B) WebSocket 수신기. 사용법:: ws = KiwoomWebSocketPriceCache(app_key, app_secret, is_mock=False) ws.start() ws.subscribe("005930") data = ws.get_price("005930") # KIS 와 동일 dict 포맷 ws.stop() ``get_price()`` 가 None 이면 → 캐시 없음/만료 → 호출자가 다른 소스 사용. """ # 키움 0B FID FID_PRICE = "10" # 현재가 (체결가, 부호 포함) FID_CHANGE = "11" # 전일대비 FID_CHANGE_PCT = "12" # 등락률 FID_CUM_VOL = "13" # 누적거래량 FID_OPEN = "16" FID_HIGH = "17" FID_LOW = "18" FID_TICK_VOL = "15" FID_TICK_TIME = "20" # 한도 (키움 권장) MAX_SUBSCRIPTIONS_PER_GROUP = 100 # grp_no=1 그룹 1개당 GROUP_NO = "1" SUB_TYPE = "0B" # 주식체결 # 재연결 정책 RECONNECT_BASE_DELAY_SEC = 5.0 RECONNECT_MAX_DELAY_SEC = 300.0 MAX_RECONNECTS_PER_HOUR = 6 MAX_RECONNECT_ATTEMPTS = 10 STABLE_CONN_RESET_SEC = 300.0 # 5분 안정 연결 후 끊기면 카운터 초기화 # 토큰 캐시 — KiwoomTokenManager 가 알아서 처리하지만 보수적 만료 버퍼 TOKEN_REFRESH_BUFFER_SEC = 600 def __init__( self, app_key: str, app_secret: str, is_mock: bool = False, ): self.app_key = app_key self.app_secret = app_secret self.is_mock = is_mock # URL (env/DB 로 재정의 가능) _default_real = "wss://api.kiwoom.com:10000/api/dostk/websocket" _default_mock = "wss://mockapi.kiwoom.com:10000/api/dostk/websocket" self._ws_url = ( get_env_from_db("KIWOOM_WS_URL_MOCK", _default_mock) if is_mock else get_env_from_db("KIWOOM_WS_URL_REAL", _default_real) ) # 메모리 캐시 — KIS WS 와 동일 포맷 (data + ts) self._cache: Dict[str, Dict] = {} self._cache_lock = threading.Lock() # 구독 종목 self._subscribed: Set[str] = set() self._sub_lock = threading.Lock() # 연결 상태 self._ws = None self._ws_thread: Optional[threading.Thread] = None self._running = False self._connected = False self._authenticated = False # LOGIN 응답 OK 받기 전엔 REG 못 보냄 # 재연결 추적 self._reconnect_count = 0 self._reconnect_times: list = [] self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC self._last_connect_time: float = 0.0 # websocket-client lib try: import websocket as _ws_lib # type: ignore self._ws_lib = _ws_lib self._available = True except ImportError: self._ws_lib = None self._available = False logger.warning("⚠️ websocket-client 미설치 — 키움 WS 사용 불가") # ------------------------------------------------------------------ # 외부 API # ------------------------------------------------------------------ def start(self) -> bool: """백그라운드 수신 스레드 기동.""" if not self._available: logger.warning("키움 WS 라이브러리 없음 → start 무시") return False if self._running: return True if not self.app_key or not self.app_secret: logger.warning("⚠️ 키움 키 없음 → 키움 WS 비활성") return False self._running = True self._ws_thread = threading.Thread( target=self._run_loop, daemon=True, name="KiwoomWS", ) self._ws_thread.start() logger.info("✅ 키움 WebSocket 수신 스레드 시작 (mock=%s, url=%s)", self.is_mock, self._ws_url) return True def stop(self) -> None: """수신 스레드 종료 + 소켓 닫기.""" self._running = False try: if self._ws is not None: self._ws.close() except Exception: pass def subscribe(self, code: str) -> bool: """단일 종목 등록. WS 연결되어 있어야 즉시 발송, 아니면 대기 큐에만 추가.""" code = (code or "").strip() if not code: return False with self._sub_lock: if code in self._subscribed: return True if len(self._subscribed) >= self.MAX_SUBSCRIPTIONS_PER_GROUP: logger.warning( "⚠️ 키움 WS 구독 한도 초과 (%d/%d) — %s 등록 거절", len(self._subscribed), self.MAX_SUBSCRIPTIONS_PER_GROUP, code, ) return False self._subscribed.add(code) # 연결되어 있을 때만 즉시 발송. 아니면 _on_open 에서 일괄 발송. if self._connected and self._authenticated: return self._send_reg([code]) return True def unsubscribe(self, code: str) -> bool: """단일 종목 등록 해지.""" code = (code or "").strip() with self._sub_lock: if code not in self._subscribed: return False self._subscribed.discard(code) if self._connected and self._authenticated: return self._send_remove([code]) return True def get_price(self, code: str, max_age_sec: float = 5.0) -> Optional[Dict]: """KIS WS ``get_price`` 와 동일 포맷 반환. 반환:: { "stck_prpr": "73900", # 현재가 "stck_oprc": "73000", # 시가 "stck_hgpr": "74500", # 고가 "stck_lwpr": "72800", # 저가 "prdy_vrss": "200", # 전일 대비 "prdy_ctrt": "0.27", # 등락률 "_age_ms": 123, # 캐시 나이 (ms) — 검증용 메타 } ``max_age_sec`` 초 초과면 None. """ with self._cache_lock: entry = self._cache.get(code) if not entry: return None age_sec = time.time() - entry.get("ts", 0) if age_sec > max_age_sec: return None data = dict(entry["data"]) data["_age_ms"] = int(age_sec * 1000) return data def is_connected(self) -> bool: return bool(self._connected and self._authenticated) def subscribed_count(self) -> int: with self._sub_lock: return len(self._subscribed) # ------------------------------------------------------------------ # 내부: 메인 수신 루프 # ------------------------------------------------------------------ def _run_loop(self) -> None: while self._running: try: self._connect_and_serve() except Exception as e: logger.warning("키움 WS 루프 예외: %s", e) finally: self._connected = False self._authenticated = False if self._running: self._reconnect_with_backoff() def _connect_and_serve(self) -> None: """단일 연결 수명. blocking. 끊기면 반환 → 호출자가 backoff 후 재호출.""" if not self._ws_lib: return # 토큰 발급 token = self._get_kiwoom_token() if not token: logger.warning("⚠️ 키움 토큰 발급 실패 → WS 연결 보류 (60s)") time.sleep(60) return # WebSocketApp 생성 self._ws = self._ws_lib.WebSocketApp( self._ws_url, on_open=self._on_open(token), on_message=self._on_message, on_error=self._on_error, on_close=self._on_close, ) self._last_connect_time = time.time() # blocking — 연결 종료까지 여기서 대기 self._ws.run_forever(ping_interval=30, ping_timeout=10) def _on_open(self, token: str): """on_open 콜백 팩토리 — token 캡처 후 LOGIN 발송.""" def _handler(ws): self._connected = True self._authenticated = False try: ws.send(json.dumps({"trnm": "LOGIN", "token": token})) logger.info("📡 키움 WS 연결 → LOGIN 발송") except Exception as e: logger.warning("키움 WS LOGIN 발송 실패: %s", e) return _handler def _on_message(self, ws, message: str) -> None: """수신 메시지 디스패치 (LOGIN ack / REG ack / REAL / PING).""" try: msg = json.loads(message) except Exception: return trnm = msg.get("trnm", "") if trnm == "PING": # 키움 PING → 그대로 echo (서버 정책) try: ws.send(message) except Exception: pass return if trnm == "LOGIN": rc = msg.get("return_code") rm = msg.get("return_msg", "") if rc == 0: self._authenticated = True logger.info("✅ 키움 WS LOGIN OK") # 누적된 구독 일괄 등록 with self._sub_lock: pending = list(self._subscribed) if pending: self._send_reg(pending) # 안정 연결 카운터 초기화 (지속 5분 이상 연결됐다면) # 여기선 LOGIN 직후라 의미 없음, _periodic_reset_ok() 에서 처리 else: logger.warning("❌ 키움 WS LOGIN 실패 rc=%s msg=%s", rc, rm) try: ws.close() except Exception: pass return if trnm == "REG": rc = msg.get("return_code") if rc != 0: logger.warning("⚠️ 키움 WS REG 실패: %s", msg.get("return_msg", "")) return if trnm == "REMOVE": return # 무관심 if trnm == "REAL": self._handle_real(msg) return def _handle_real(self, msg: dict) -> None: """실시간 데이터 처리 — 0B 만 사용 (확장 가능).""" items = msg.get("data") or [] for item in items: sub_type = str(item.get("type", "")).strip() if sub_type != self.SUB_TYPE: continue code = str(item.get("item", "")).strip() values = item.get("values") or {} self._cache_tick(code, values) def _cache_tick(self, code: str, values: dict) -> None: """0B 체결 → 메모리 캐시 갱신 (KIS 와 동일 포맷).""" if not code or not values: return try: # 키움은 가격 부호로 등락 표시 — 절대값 취함 price_raw = values.get(self.FID_PRICE, "0") price = abs(float(str(price_raw).replace(",", ""))) if price <= 0: return def _abs_str(v: str) -> str: try: return str(int(abs(float(str(v).replace(",", ""))))) except (ValueError, TypeError): return "0" # KIS inquire_price 호환 필드 (스칼라 dict) data_compat = { "stck_prpr": str(int(price)), "stck_oprc": _abs_str(values.get(self.FID_OPEN, "0")), "stck_hgpr": _abs_str(values.get(self.FID_HIGH, "0")), "stck_lwpr": _abs_str(values.get(self.FID_LOW, "0")), "prdy_vrss": _abs_str(values.get(self.FID_CHANGE, "0")), "prdy_ctrt": str(values.get(self.FID_CHANGE_PCT, "0")), } with self._cache_lock: self._cache[code] = {"data": data_compat, "ts": time.time()} except (ValueError, TypeError) as e: logger.debug("키움 0B 파싱 오류 %s: %s", code, e) def _on_error(self, ws, error) -> None: self._connected = False self._authenticated = False logger.warning("⚠️ 키움 WS 오류: %s", error) def _on_close(self, ws, close_status_code, close_msg) -> None: self._connected = False self._authenticated = False logger.info("🔌 키움 WS 연결 종료 (code=%s msg=%s)", close_status_code, close_msg or "") # ------------------------------------------------------------------ # 내부: 등록/해지 메시지 발송 # ------------------------------------------------------------------ def _send_reg(self, codes: list) -> bool: """REG 발송 (그룹 1, 0B 타입). 한 번에 여러 종목 OK.""" if not codes or not self._ws: return False try: self._ws.send(json.dumps({ "trnm": "REG", "grp_no": self.GROUP_NO, "refresh": "1", # 재시작 시 등록 유지 "data": [{"item": list(codes), "type": [self.SUB_TYPE]}], })) logger.info("📡 키움 WS REG 발송: %d종목 (총 %d/%d)", len(codes), self.subscribed_count(), self.MAX_SUBSCRIPTIONS_PER_GROUP) return True except Exception as e: logger.warning("키움 WS REG 실패: %s", e) return False def _send_remove(self, codes: list) -> bool: """REMOVE 발송.""" if not codes or not self._ws: return False try: self._ws.send(json.dumps({ "trnm": "REMOVE", "grp_no": self.GROUP_NO, "data": [{"item": list(codes), "type": [self.SUB_TYPE]}], })) return True except Exception as e: logger.warning("키움 WS REMOVE 실패: %s", e) return False # ------------------------------------------------------------------ # 내부: 토큰 # ------------------------------------------------------------------ def _get_kiwoom_token(self) -> Optional[str]: """``kis_ws.KiwoomTokenManager`` 싱글톤 풀 재사용.""" try: from kis_ws import _get_kiwoom_token_cached # type: ignore except ImportError: logger.warning("kis_ws._get_kiwoom_token_cached import 실패") return None return _get_kiwoom_token_cached(self.app_key, self.app_secret, self.is_mock) # ------------------------------------------------------------------ # 내부: 재연결 백오프 (KIS WS 와 동일 정책) # ------------------------------------------------------------------ def _reconnect_with_backoff(self) -> None: now = time.time() # 5분 이상 안정 연결 후 끊긴 거면 카운터 초기화 (정상 재연결 케이스) if (self._last_connect_time > 0 and now - self._last_connect_time >= self.STABLE_CONN_RESET_SEC): self._reconnect_count = 0 self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC # 시간당 한도 초과 → 1시간 대기 cutoff = now - 3600 self._reconnect_times = [t for t in self._reconnect_times if t > cutoff] if len(self._reconnect_times) >= self.MAX_RECONNECTS_PER_HOUR: wait = 3600 - (now - self._reconnect_times[0]) logger.warning("⚠️ 키움 WS 시간당 재연결 한도 초과 → %ds 대기", int(wait)) time.sleep(max(60.0, wait)) return # 총 한도 초과 → 비활성 if self._reconnect_count >= self.MAX_RECONNECT_ATTEMPTS: logger.warning("⛔ 키움 WS 총 재연결 한도 초과 → 자동 비활성") self._running = False return delay = min(self._reconnect_delay, self.RECONNECT_MAX_DELAY_SEC) logger.info("⏳ 키움 WS %ds 후 재연결 시도 (#%d)", int(delay), self._reconnect_count + 1) time.sleep(delay) self._reconnect_times.append(time.time()) self._reconnect_count += 1 self._reconnect_delay = min(self._reconnect_delay * 2, self.RECONNECT_MAX_DELAY_SEC)