From 3fa9eb9bf7026f898422e5d8b134f8c6a28a5de3 Mon Sep 17 00:00:00 2001 From: Hwang Date: Tue, 5 May 2026 21:28:22 +0900 Subject: [PATCH] =?UTF-8?q?feat(ws):=20=ED=82=A4=EC=9B=80=20WS=20=EC=8B=9C?= =?UTF-8?q?=EC=84=B8=20=EB=A7=88=EC=9D=B4=EA=B7=B8=EB=A0=88=EC=9D=B4?= =?UTF-8?q?=EC=85=98=20=EA=B2=80=EC=A6=9D=20=EC=9D=B8=ED=94=84=EB=9D=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit KIS WS(41 한도) → 키움 WS(100 한도) 본격 전환 전, 두 소스를 동시 운영해 가격 일치성을 데이터로 검증하기 위한 인프라. 신규 - kiwoom_ws.py: 키움 WS 클라이언트 (KIS WS 와 동일 get_price 인터페이스, 메모리 dict 캐시) - kis_trader/network/ws_validator.py: 5초마다 KIS↔키움 가격 비교, ws_price_validation 테이블에 1행 INSERT, |diff|≥WARN_PCT 시 WARN 로그 - test_kiwoom_ws.py: 키움 WS 단독 동작 확인 스크립트 (토큰/LOGIN/REG/시세) 수정 - database.py: ws_price_validation 테이블 + ENV 키 (WS_PROVIDER, WS_VALIDATION_INTERVAL_SEC, WS_VALIDATION_DIFF_WARN_PCT) + insert_ws_price_validation/get_ws_validation_stats 헬퍼 - kis_trader/main.py: WS_PROVIDER=kis_with_validation 시 키움 WS + Validator 백그라운드 기동, 종료 시 정리 운영 영향: 0. 매매·시세 의사결정은 항상 KIS WS만 사용. 키움 WS는 백그라운드 비교 기록만. 적용: DB env_config 에 WS_PROVIDER=kis_with_validation INSERT 후 봇 재시작. Co-authored-by: Cursor --- database.py | 116 ++++++- kis_trader/main.py | 87 +++++ kis_trader/network/ws_validator.py | 213 ++++++++++++ kiwoom_ws.py | 502 +++++++++++++++++++++++++++++ test_kiwoom_ws.py | 170 ++++++++++ 5 files changed, 1087 insertions(+), 1 deletion(-) create mode 100644 kis_trader/network/ws_validator.py create mode 100644 kiwoom_ws.py create mode 100644 test_kiwoom_ws.py diff --git a/database.py b/database.py index a30b174..d89d317 100644 --- a/database.py +++ b/database.py @@ -13,7 +13,7 @@ import os import datetime import logging import threading -from typing import Dict, List, Optional, Tuple +from typing import Any, Dict, List, Optional, Tuple try: import pymysql @@ -417,6 +417,17 @@ ENV_CONFIG_KEYS = ( # 미설정(0/빈값) 시 글로벌 MAX_STOCKS 로 폴백 → 구버전 호환. # 권장: 합계 ≤ MAX_STOCKS (계좌 슬롯 분산), 예: 3+2+2=7. "SCALP_MAX_STOCKS", "SHORT_MAX_STOCKS", "BREAKOUT_MAX_STOCKS", + # ── 시세 WS 공급자 토글 (키움 시세 마이그레이션) ───────────────── + # 운영(매매 의사결정)에는 항상 KIS WS 만 사용. 키움 WS 는 검증 모드에서만 + # 백그라운드 동시 구독 → ws_price_validation 테이블에 가격 비교 기록. + # kis_only : 현행 (기본). 키움 WS 미기동. + # kis_with_validation : KIS WS 운영 + 키움 WS 검증 동시 (매매 영향 없음) + # kiwoom_only : 시세를 키움으로 전환 (검증 통과 후에만 사용) + "WS_PROVIDER", + # 검증 비교 주기(초) — 너무 짧으면 부하, 너무 길면 표본 부족. 기본 5. + "WS_VALIDATION_INTERVAL_SEC", + # 차이 경고 임계(%). |diff| 가 이 값 이상이면 WARN 로그. 기본 0.10. + "WS_VALIDATION_DIFF_WARN_PCT", ) @@ -762,6 +773,30 @@ class TradeDB: logger.info("📌 target_candidates_history 테이블 확인/생성") except Exception as e: logger.warning(f"migrate target_candidates_history 실패(이력 미적재 가능): {e}") + # ── ws_price_validation (KIS↔키움 시세 검증, 마이그레이션 단계용) ──── + # 5초마다 같은 종목의 KIS WS 가격과 키움 WS 가격을 비교해 한 행 INSERT. + # diff_pct = (kiwoom - kis) / kis × 100. + # 운영에는 영향 없음 (검증 모드 ON 일 때만 채워짐). 1~2주 누적 후 + # 통계 분석 → 본격 마이그레이션 결정 근거. + try: + self.conn.execute(""" + CREATE TABLE IF NOT EXISTS ws_price_validation ( + id BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY, + ts DATETIME(3) NOT NULL, + code VARCHAR(20) NOT NULL, + kis_price DOUBLE, + kiwoom_price DOUBLE, + diff_pct DOUBLE, + kis_age_ms INT, + kiwoom_age_ms INT, + INDEX idx_ts (ts), + INDEX idx_code (code), + INDEX idx_diff (diff_pct) + ) CHARACTER SET utf8mb4 + """) + logger.info("📌 ws_price_validation 테이블 확인/생성") + except Exception as e: + logger.warning(f"migrate ws_price_validation 실패: {e}") def _migrate_env_config_to_columns(self): """env_config가 예전 JSON 컬럼(snapshot_json)이면 컬럼 스키마로 이전""" @@ -1837,6 +1872,85 @@ class TradeDB: logger.error(f"❌ 날짜별 조회 실패: {e}") return [] + # ============================================================ + # [ws_price_validation] KIS↔키움 시세 비교 검증 + # ============================================================ + + def insert_ws_price_validation( + self, + *, + code: str, + kis_price: Optional[float], + kiwoom_price: Optional[float], + kis_age_ms: Optional[int] = None, + kiwoom_age_ms: Optional[int] = None, + ) -> bool: + """단일 비교 결과 1행 INSERT. + + 둘 다 None 이면 저장 안 함. 한쪽만 있어도 저장(소스별 가용성 분석용). + diff_pct 는 둘 다 있을 때만 계산. + """ + if kis_price is None and kiwoom_price is None: + return False + diff_pct: Optional[float] = None + if kis_price not in (None, 0) and kiwoom_price is not None: + try: + diff_pct = (float(kiwoom_price) - float(kis_price)) / float(kis_price) * 100.0 + except (ValueError, ZeroDivisionError): + diff_pct = None + try: + now = datetime.datetime.now() + self.conn.execute( + "INSERT INTO ws_price_validation " + "(ts, code, kis_price, kiwoom_price, diff_pct, kis_age_ms, kiwoom_age_ms) " + "VALUES (%s, %s, %s, %s, %s, %s, %s)", + (now, code, kis_price, kiwoom_price, diff_pct, kis_age_ms, kiwoom_age_ms), + ) + return True + except Exception as e: + logger.debug("ws_price_validation INSERT 실패: %s", e) + return False + + def get_ws_validation_stats( + self, *, hours: int = 24, code: Optional[str] = None, + ) -> Dict[str, Any]: + """최근 N시간 검증 통계 (운영자용 분석). + + Returns: + { + "samples": 1234, + "both_present": 1100, # KIS·키움 둘 다 가격 있던 비율 + "avg_diff_pct": 0.012, + "max_abs_diff_pct": 0.45, + "stddev_diff_pct": 0.08, + "kis_only": 80, # KIS 만 가격 있던 횟수 (키움 미수신) + "kiwoom_only": 30, # 키움 만 가격 있던 횟수 + } + """ + try: + where = ["ts >= NOW() - INTERVAL %s HOUR"] + args: List[Any] = [hours] + if code: + where.append("code = %s") + args.append(code) + wsql = " AND ".join(where) + row = self.conn.execute(f""" + SELECT + COUNT(*) AS samples, + SUM(kis_price IS NOT NULL AND kiwoom_price IS NOT NULL) AS both_present, + AVG(diff_pct) AS avg_diff_pct, + MAX(ABS(diff_pct)) AS max_abs_diff_pct, + STDDEV(diff_pct) AS stddev_diff_pct, + SUM(kis_price IS NOT NULL AND kiwoom_price IS NULL) AS kis_only, + SUM(kis_price IS NULL AND kiwoom_price IS NOT NULL) AS kiwoom_only + FROM ws_price_validation + WHERE {wsql} + """, args).fetchone() + return dict(row) if row else {} + except Exception as e: + logger.debug("ws_validation_stats 조회 실패: %s", e) + return {} + # ============================================================ # [env_config] 관리자용 env (INSERT만, 최신 1건 = 살아있는 값) # ============================================================ diff --git a/kis_trader/main.py b/kis_trader/main.py index 6b82031..322efbd 100644 --- a/kis_trader/main.py +++ b/kis_trader/main.py @@ -111,6 +111,9 @@ class TradingOrchestrator: # 시장 급락 서킷브레이커 (KOSPI/KOSDAQ 지수 폭락 시 신규 매수 전면 차단). # 지수 조회는 모의 도메인 미지원 가능성 → 실키 우선, 폴백 self.client. self.market_guard: MarketGuard | None = None + # 시세 마이그레이션 검증 (WS_PROVIDER=kis_with_validation 시만 기동) + self.kiwoom_ws = None + self.ws_validator = None self.strategies: List[BaseStrategy] = [] self._stop = False @@ -221,6 +224,9 @@ class TradingOrchestrator: # 시장 급락 서킷브레이커 매니저 기동 (활성/비활성 무관, 토글로 즉시 전환) self._start_market_guard() + # 시세 WS 마이그레이션 검증 (WS_PROVIDER 토글에 따라 키움 WS 동시 기동) + self._start_ws_validator() + # 전략 등록 후 쓰레드 시작 (각자 독립 쓰레드) self._register_strategies() for strat in self.strategies: @@ -400,6 +406,77 @@ class TradingOrchestrator: logger.error("MarketGuard 기동 실패 (전략은 가드 없이 동작): %s", e) self.market_guard = None + # ------------------------------------------------------------------ + def _start_ws_validator(self) -> None: + """시세 WS 마이그레이션 검증 — 키움 WS + Validator 백그라운드 기동. + + ``WS_PROVIDER`` env 토글: + kis_only : 미기동 (기본). 운영 100% 동일. + kis_with_validation : KIS WS 운영 + 키움 WS 검증 동시 (매매 영향 없음). + kiwoom_only : 추후 단계 (검증 통과 후 시세 키움 전환). + + 매매·시세 의사결정은 항상 KIS WS 만 사용. 키움 WS 는 ``ws_price_validation`` + 테이블에 기록만 함. 1~2주 누적 후 통계 분석으로 마이그레이션 결정. + """ + from .utils.env import get_env_from_db + + provider = (get_env_from_db("WS_PROVIDER", "kis_only") or "kis_only").strip().lower() + if provider == "kis_only": + logger.info("ℹ️ WS_PROVIDER=kis_only → 키움 WS 검증 비활성") + return + if provider not in ("kis_with_validation",): + logger.warning( + "⚠️ WS_PROVIDER='%s' 알 수 없음 → kis_only 로 폴백 (검증 비활성)", + provider, + ) + return + + # 키움 키 로드 (kis_ws._get_kiwoom_creds 재사용) + try: + from kis_ws import _get_kiwoom_creds # type: ignore + app_key, app_secret, is_mock = _get_kiwoom_creds(self.db) + except Exception as e: + logger.warning("키움 자격증명 로드 실패 → 검증 비활성: %s", e) + return + if not app_key or not app_secret: + logger.warning("키움 키 미설정 → WS 검증 비활성") + return + + # 키움 WS 인스턴스 + try: + from kiwoom_ws import KiwoomWebSocketPriceCache # type: ignore + self.kiwoom_ws = KiwoomWebSocketPriceCache( + app_key, app_secret, is_mock=is_mock, + ) + if not self.kiwoom_ws.start(): + logger.warning("키움 WS 시작 실패 → 검증 비활성") + self.kiwoom_ws = None + return + except Exception as e: + logger.warning("키움 WS 인스턴스 생성 실패: %s", e) + self.kiwoom_ws = None + return + + # KIS WS 핸들 가져오기 (WSManager.ws_cache) + kis_ws_handle = getattr(self.ws, "ws_cache", None) + if kis_ws_handle is None: + logger.warning("KIS WS 핸들 미발견 → Validator 비활성 (검증 불가)") + return + + # Validator 기동 + try: + from .network.ws_validator import WSPriceValidator + self.ws_validator = WSPriceValidator( + kis_ws=kis_ws_handle, + kiwoom_ws=self.kiwoom_ws, + db=self.db, + ) + self.ws_validator.start() + logger.info("🔬 [시세 검증 모드] WS_PROVIDER=%s — KIS↔키움 동시 운영", provider) + except Exception as e: + logger.warning("Validator 기동 실패: %s", e) + self.ws_validator = None + def _build_market_client(self) -> KISClient: """ 시세/조회 전용 실전 KISClient 생성 (init 단계). @@ -816,6 +893,16 @@ class TradingOrchestrator: self.ranking_mgr.stop() except Exception: pass + try: + if self.ws_validator: + self.ws_validator.stop() + except Exception: + pass + try: + if self.kiwoom_ws: + self.kiwoom_ws.stop() + except Exception: + pass try: self.ws.stop() except Exception: diff --git a/kis_trader/network/ws_validator.py b/kis_trader/network/ws_validator.py new file mode 100644 index 0000000..3e630c9 --- /dev/null +++ b/kis_trader/network/ws_validator.py @@ -0,0 +1,213 @@ +""" +kis_trader/network/ws_validator.py — KIS↔키움 WS 가격 검증 +============================================================ +주기적으로 같은 종목의 KIS WS 캐시와 키움 WS 캐시를 비교해 +``ws_price_validation`` 테이블에 1행 INSERT. + +운영(매매)에는 영향 없음 — **읽기만** 함. +검증 모드(``WS_PROVIDER=kis_with_validation``)에서만 기동. + +데이터 흐름:: + + [KIS WS] ─┐ + │ 각각 메모리 dict 캐시 + [키움 WS] ─┘ + │ + ▼ + [Validator] ── 5s 주기 ──► ws_price_validation 테이블 + │ │ + │ └─→ |diff| ≥ WARN_PCT 면 WARN 로그 + │ + └─→ 24h/1주 통계 분석 → 마이그레이션 전환 결정 근거 + +비교 대상 종목 +-------------- +KIS WS 가 구독 중인 종목 = 봇이 매매에 쓰는 가격이 실제로 들어오고 있는 종목. +키움 WS 도 같은 종목을 구독하도록 동기화 (subscribe). + +env_config 토글 +--------------- +``WS_VALIDATION_INTERVAL_SEC`` (기본 5) +``WS_VALIDATION_DIFF_WARN_PCT`` (기본 0.10 — 0.1%p) +""" +from __future__ import annotations + +import threading +import time +from typing import Optional, Set + +from ..utils.env import get_env_float, get_env_int +from ..utils.logger import get_logger + +logger = get_logger("kis_trader.ws_validator") + + +class WSPriceValidator: + """KIS↔키움 WS 가격 비교 백그라운드 워커. + + 매 N초마다: + 1) KIS WS 의 구독 중인 종목 목록을 읽음 + 2) 키움 WS 가 같은 종목을 구독하도록 동기화 + 3) 두 캐시에서 가격 조회 → DB INSERT + 4) |diff_pct| ≥ warn_pct 시 WARN 로그 + """ + + def __init__( + self, + *, + kis_ws, # kis_ws.KISWebSocketPriceCache 인스턴스 + kiwoom_ws, # kiwoom_ws.KiwoomWebSocketPriceCache 인스턴스 + db, # database.TradeDB + ): + self.kis_ws = kis_ws + self.kiwoom_ws = kiwoom_ws + self.db = db + + self._thread: Optional[threading.Thread] = None + self._running = False + self._last_warn_ts: dict = {} # code → 최근 WARN 로그 시각 (스팸 방지) + + # ------------------------------------------------------------------ + def start(self) -> bool: + if self._thread and self._thread.is_alive(): + return True + self._running = True + self._thread = threading.Thread( + target=self._loop, daemon=True, name="WSPriceValidator", + ) + self._thread.start() + logger.info( + "✅ WS 가격 검증기 시작 — KIS↔키움 비교 (interval=%ds, warn≥%.2f%%)", + self._interval_sec(), self._warn_pct(), + ) + return True + + def stop(self) -> None: + self._running = False + + # ------------------------------------------------------------------ + def _interval_sec(self) -> int: + return max(1, get_env_int("WS_VALIDATION_INTERVAL_SEC", 5)) + + def _warn_pct(self) -> float: + return max(0.0, get_env_float("WS_VALIDATION_DIFF_WARN_PCT", 0.10)) + + # ------------------------------------------------------------------ + def _kis_subscribed(self) -> Set[str]: + """KIS WS 가 현재 구독 중인 종목 set.""" + # KISWebSocketPriceCache 의 _subscribed 직접 참조 (kis_ws.py 정의) + try: + with self.kis_ws._sub_lock: # type: ignore[attr-defined] + return set(self.kis_ws._subscribed) # type: ignore[attr-defined] + except AttributeError: + return set() + + def _sync_kiwoom_subscriptions(self, target: Set[str]) -> None: + """키움 WS 구독 = KIS WS 구독 으로 맞춤.""" + try: + current = set() + with self.kiwoom_ws._sub_lock: # type: ignore[attr-defined] + current = set(self.kiwoom_ws._subscribed) # type: ignore[attr-defined] + + to_add = target - current + to_remove = current - target + for code in to_add: + self.kiwoom_ws.subscribe(code) + for code in to_remove: + self.kiwoom_ws.unsubscribe(code) + except Exception as e: + logger.debug("키움 구독 동기화 실패: %s", e) + + # ------------------------------------------------------------------ + def _loop(self) -> None: + # 시작 직후 KIS/키움 둘 다 캐시 채워질 시간 약간 줌 + time.sleep(15) + while self._running: + try: + self._tick() + except Exception as e: + logger.warning("WS 검증기 tick 예외: %s", e) + time.sleep(self._interval_sec()) + + def _tick(self) -> None: + """1회 비교.""" + codes = self._kis_subscribed() + if not codes: + return + + # 키움 WS 구독 동기화 (구독 안 된 종목은 데이터 없으니 의미 없음) + self._sync_kiwoom_subscriptions(codes) + + # 키움 WS 연결되어 있어야 의미 있음. 미연결이면 한쪽만 기록. + kiwoom_ready = self.kiwoom_ws.is_connected() + + warn_pct = self._warn_pct() + warn_count = 0 + sample_count = 0 + + for code in codes: + kis_data = self.kis_ws.get_price(code, max_age_sec=10.0) + kw_data = self.kiwoom_ws.get_price(code, max_age_sec=10.0) if kiwoom_ready else None + + kis_price = self._parse_price(kis_data, "stck_prpr") + kw_price = self._parse_price(kw_data, "stck_prpr") + + kis_age = self._parse_age_ms(kis_data) + kw_age = self._parse_age_ms(kw_data) + + if kis_price is None and kw_price is None: + continue + + # DB 1행 INSERT + self.db.insert_ws_price_validation( + code=code, + kis_price=kis_price, + kiwoom_price=kw_price, + kis_age_ms=kis_age, + kiwoom_age_ms=kw_age, + ) + sample_count += 1 + + # 차이 경고 + if kis_price not in (None, 0) and kw_price is not None: + diff_pct = (kw_price - kis_price) / kis_price * 100.0 + if abs(diff_pct) >= warn_pct: + # 종목당 60초 1회만 로그 (스팸 방지) + now = time.time() + last = self._last_warn_ts.get(code, 0) + if now - last >= 60: + self._last_warn_ts[code] = now + logger.warning( + "⚠️ [WS 검증] %s 가격 차이 %.3f%% (KIS=%.0f, 키움=%.0f)", + code, diff_pct, kis_price, kw_price, + ) + warn_count += 1 + + if sample_count > 0: + logger.debug( + "📊 [WS 검증] tick: %d종목 비교 (warn=%d, 키움연결=%s)", + sample_count, warn_count, kiwoom_ready, + ) + + # ------------------------------------------------------------------ + @staticmethod + def _parse_price(data: Optional[dict], key: str) -> Optional[float]: + if not data: + return None + try: + v = float(str(data.get(key, "0")).replace(",", "")) + return abs(v) if v else None + except (ValueError, TypeError): + return None + + @staticmethod + def _parse_age_ms(data: Optional[dict]) -> Optional[int]: + if not data: + return None + v = data.get("_age_ms") + if v is None: + return None + try: + return int(v) + except (ValueError, TypeError): + return None diff --git a/kiwoom_ws.py b/kiwoom_ws.py new file mode 100644 index 0000000..ac8af70 --- /dev/null +++ b/kiwoom_ws.py @@ -0,0 +1,502 @@ +""" +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) diff --git a/test_kiwoom_ws.py b/test_kiwoom_ws.py new file mode 100644 index 0000000..6e2dbcc --- /dev/null +++ b/test_kiwoom_ws.py @@ -0,0 +1,170 @@ +#!/usr/bin/env python3 +""" +test_kiwoom_ws.py — 키움 WebSocket 시세 단독 검증 스크립트 +========================================================= + +봇 안 띄우고 ``kiwoom_ws.KiwoomWebSocketPriceCache`` 만 직접 띄워서 +연결·로그인·등록·시세 수신이 정상인지 30~60초 안에 확인. + +사용법 +------ + +기본 (DB 키 자동 로드, 삼성전자 + SK하이닉스 30초 구독):: + + python3 test_kiwoom_ws.py + +종목 직접 지정:: + + python3 test_kiwoom_ws.py --codes 005930,000660,005380 --duration 60 + +장외에 띄워도 LOGIN/REG ack 까지 확인 가능 (시세는 안 오지만 인증 동작 검증). + +출력 해석 +--------- + +✅ LOGIN OK → 키움 토큰·키 정상 +✅ REG OK → 종목 등록 성공 +📈 005930 73900... → 1초 1줄, 실시간 가격 수신 정상 +⚠️ 가격 수신 0건 → 장외 또는 키움 서버 정책 +❌ LOGIN 실패 → 키 / OpenAPI 신청 / 도메인 문제 (test_kiwoom_token.py 로 추가 진단) +""" +from __future__ import annotations + +import argparse +import logging +import sys +import time +from pathlib import Path + +SCRIPT_DIR = Path(__file__).resolve().parent +sys.path.insert(0, str(SCRIPT_DIR)) + + +def setup_logging(verbose: bool) -> None: + fmt = "[%(asctime)s] [%(name)s] %(message)s" + logging.basicConfig( + level=logging.DEBUG if verbose else logging.INFO, + format=fmt, + datefmt="%H:%M:%S", + ) + + +def load_kiwoom_creds() -> tuple[str, str, bool]: + """DB env_config 에서 키움 키 로드 (KIS_MOCK 따라 MOCK/REAL 자동 선택).""" + from kis_ws import _get_kiwoom_creds # type: ignore + from database import TradeDB + + db = TradeDB() + try: + return _get_kiwoom_creds(db) + finally: + db.close() + + +def main() -> int: + p = argparse.ArgumentParser(description="키움 WebSocket 시세 단독 검증") + p.add_argument( + "--codes", + default="005930,000660", + help="콤마 구분 종목코드 (기본: 삼성전자, SK하이닉스)", + ) + p.add_argument( + "--duration", type=int, default=30, + help="구독 유지 시간(초). 기본 30", + ) + p.add_argument("-v", "--verbose", action="store_true", help="DEBUG 로그") + args = p.parse_args() + + setup_logging(args.verbose) + logger = logging.getLogger("test_kiwoom_ws") + + codes = [c.strip() for c in args.codes.split(",") if c.strip()] + if not codes: + print("❌ --codes 비어있음") + return 2 + + print("=" * 70) + print("🔍 키움 WS 단독 검증") + print("=" * 70) + + # 키 로드 + try: + app_key, app_secret, is_mock = load_kiwoom_creds() + except Exception as e: + print(f"❌ 키 로드 실패: {e}") + return 1 + if not app_key or not app_secret: + print("❌ 키움 키 미설정 (DB env_config 의 KIWOOM_APP_KEY_* 확인)") + return 1 + print(f" 키움 키: {app_key[:8]}…{app_key[-4:]} is_mock={is_mock}") + print(f" 대상 종목: {codes}") + print(f" 유지 시간: {args.duration}초") + print() + + # WS 인스턴스 + from kiwoom_ws import KiwoomWebSocketPriceCache + + ws = KiwoomWebSocketPriceCache(app_key, app_secret, is_mock=is_mock) + if not ws.start(): + print("❌ 키움 WS start() 실패") + return 1 + + # 종목 구독 + for code in codes: + ws.subscribe(code) + + # 연결 + LOGIN 대기 (최대 15초) + deadline = time.time() + 15 + while time.time() < deadline: + if ws.is_connected(): + print(f"✅ LOGIN OK ({int(time.time() - (deadline - 15))}s)") + break + time.sleep(0.5) + else: + print("⚠️ LOGIN 미완료 (15초 timeout). 키움 토큰/네트워크 확인 필요") + ws.stop() + return 1 + + # 가격 수신 모니터링 + print() + print(f"📡 {args.duration}초간 시세 모니터링 (1초 간격 출력)...") + end_ts = time.time() + args.duration + last_dump_ts = 0.0 + received: dict = {} + while time.time() < end_ts: + now = time.time() + if now - last_dump_ts >= 1.0: + last_dump_ts = now + for code in codes: + d = ws.get_price(code, max_age_sec=10.0) + if d: + px = d.get("stck_prpr", "?") + chg = d.get("prdy_ctrt", "?") + age = d.get("_age_ms", "?") + print(f" 📈 {code} price={px}원 chg={chg}% age={age}ms") + received[code] = received.get(code, 0) + 1 + else: + print(f" ⏸ {code} (no data — 장외/미수신)") + time.sleep(0.2) + + # 종료 + print() + print("─" * 70) + print("종료. 수신 통계:") + for code in codes: + n = received.get(code, 0) + mark = "✅" if n > 0 else "⚠️" + print(f" {mark} {code}: {n}회 수신") + print("─" * 70) + print() + print("팁:") + print(" • 모든 종목 0회면 장외 시간이거나 키움 정책 문제일 수 있음.") + print(" • 장중에 0회면 → 키움 OpenAPI+ 신청 옵션 (실시간 시세 권한) 확인.") + print(" • LOGIN 실패면 → test_kiwoom_token.py 로 토큰 발급 자체 진단.") + + ws.stop() + return 0 if any(received.values()) else 1 + + +if __name__ == "__main__": + sys.exit(main())