""" kis_trader/ws/kiwoom_ws.py — 키움 WebSocket 실시간 시세 캐시 ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ 목적 ---- 실매 키움 실시간 시세(체결 0B · 호가 0D · 프로그램 0w) **한 소켓**. 예전에는 KIS 41 한도 대비 검증용이었으나, 지금은 틱/호가 본체. ``LIVE_TICK_PROVIDER`` 는 이 모듈이 아니라 ``WSManager.get_price`` 읽기 순서. ``WS_SUBSCRIBE_KIS_MINIMAL`` 은 한투 구독 명단. 키움 후보 구독과 무관. 설계 원칙 (KIS WS 와 동일) -------------------------- - ``Dict[str, Dict]`` 메모리 캐시 (락만 잠그고 마이크로초 read/write) - ``get_price(code)`` 인터페이스를 KIS WS 와 100% 동일 포맷으로 제공 - ``KiwoomTokenManager`` 싱글톤 (kis_ws.py 안에 있음) 재사용 - 재연결 백오프, 종목 등록/해지 키움 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) 토글 (DB env_config) -------------------- 키움 WS 기동: Orchestrator ``_start_ws_validator`` — 키 있으면 기동. ``WS_PROVIDER=kis_only`` 로 끄지 않음(구식 주석). ``KIWOOM_WS_URL_REAL`` / ``KIWOOM_WS_URL_MOCK`` (URL 재정의용) ``KIWOOM_WS_REG_CHUNK_SIZE`` / ``KIWOOM_WS_REG_GAP_SEC`` / ``KIWOOM_WS_REG_DEBOUNCE_SEC`` — REG(TRNM) 초당 허용 건수 초과 방지(배치·전송 간격·단건 디바운스). ``KIWOOM_WS_REG_REFRESH`` (기본 1) — REG ``refresh``. 0=기존 그룹 구독 해지 후 재등록. LOGIN 직후 일괄 REG **첫 청크만** 0. 장중 종목 추가 REG는 항상 1 (청크마다 0이면 앞 종목이 지워짐). 장중 상시 0 비추(catch-up 재점화). """ from __future__ import annotations import json import logging import threading import time from collections import defaultdict from typing import Any, Callable, Dict, Iterable, List, Optional, Set logger = logging.getLogger("KiwoomWebSocket") from .ws_reconnect_backoff import ws_reconnect_delay_for_attempt # ── 환경변수 헬퍼 (KIS WS 와 동일 패턴, fallback 포함) ──────────────────── try: from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int # noqa: F401 except ImportError: try: from kis_long_ver1 import get_env_float, 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 def get_env_float(key, default): # type: ignore[misc] return float(default) def get_env_bool(key, default=False): # type: ignore[misc] return default try: from .orderbook_cache import OrderbookCache, OrderbookSnapshot except ImportError: OrderbookCache = None # type: ignore[misc, assignment] OrderbookSnapshot = None # type: ignore[misc, assignment] try: from .program_cache import ProgramCache, ProgramSnapshot except ImportError: ProgramCache = None # type: ignore[misc, assignment] ProgramSnapshot = None # type: ignore[misc, assignment] from .kiwoom_ws_diag import log_kiwoom_api_msg, log_kiwoom_ws # ────────────────────────────────────────────────────────────────────── # 메인 클래스 # ────────────────────────────────────────────────────────────────────── 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" FID_EXEC_STRENGTH = "228" # 체결강도 (%) — 키움 0B 공식 FID_UPPER_LIMIT_TIME = "567" # 상한가발생시간 HHmmss FID_LOWER_LIMIT_TIME = "568" # 하한가발생시간 HHmmss # 한도 (키움 권장) MAX_SUBSCRIPTIONS_PER_GROUP = 100 # grp_no=1 그룹 1개당 GROUP_NO = "1" SUB_TYPE = "0B" # 주식체결 SUB_TYPE_ORDERBOOK = "0D" # 주식호가잔량 SUB_TYPE_PROGRAM = "0w" # 종목프로그램매매 # 재연결 sleep — ws_reconnect_backoff (기본 1,3,5,7,10초) 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() # 종목별 마지막 FID20 — 재연결 catch-up(시간 후퇴/지연) vs 시계 동결 구분 self._fid20_last: Dict[str, Any] = {} # 현재 FID20 를 처음 본 벽시계 — 같은 초 연속 체결 vs 아침 동결(수 분) 구분 self._fid20_first_wall: Dict[str, float] = {} self._stale_skip_log_ts: Dict[str, float] = {} self._price_listeners: list = [] self._price_listener_lock = threading.Lock() self._orderbook_cache = OrderbookCache() if OrderbookCache else None self._program_cache = ProgramCache() if ProgramCache else None # 구독 종목 self._subscribed: Set[str] = set() self._sub_lock = threading.Lock() # 선택: KIS CandleAggregator 에 틱 전달 (후보 종목만 필터링 가능) self._candle_agg: Any = None # None = 구독 전 종목 틱을 집계기에 전달, Set = 해당 코드만 전달 self._candle_agg_codes: Optional[Set[str]] = None # TickRecorder (C안 ws_ticks) — 필터는 TickRecorder.set_record_codes() self._tick_recorder: Any = None self._trigger_snapshot_recorder: Any = None # 연결 상태 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._last_connect_time: float = 0.0 # REG 레이트리밋 회피: 단건 subscribe 는 디바운스 후 묶어 전송 self._reg_batch_codes: Set[str] = set() self._reg_timer: Optional[threading.Timer] = None self._reg_timer_lock = threading.Lock() # 조건검색 등 외부 모듈 — **동일 WS 세션 공유** (키움은 토큰당 1접속) self._ext_handler_lock = threading.Lock() self._ext_handlers: Dict[str, List[Callable]] = defaultdict(list) self._login_callbacks: List[Callable] = [] self._login_cb_lock = threading.Lock() # 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 _unsubscribe_all_before_close(self, reason: str = "", *, clear_ram: bool = True) -> int: """세션 close 전 서버 REMOVE. LOGIN 전(구독 없음)이면 0. RAM 은 ``clear_ram=True`` (stop) 일 때만 비운다. """ with self._sub_lock: codes = sorted(self._subscribed) if clear_ram: self._reg_batch_codes.clear() self._subscribed.clear() if not codes or not self._connected or not self._authenticated or self._ws is None: return 0 n_ok = 0 try: chunk = self._reg_chunk_size() gap = self._reg_gap_sec() for i in range(0, len(codes), chunk): part = codes[i:i + chunk] if self._send_remove(part): n_ok += len(part) if i + chunk < len(codes): time.sleep(gap) logger.info( "✅ 키움 WS 종료 전 REMOVE %d/%d종목 (%s)", n_ok, len(codes), reason or "-", ) except Exception as e: logger.warning("키움 WS 종료 전 REMOVE 실패: %s", e) return n_ok def stop(self) -> None: """수신 스레드 종료 + 구독 REMOVE 후 소켓 닫기 (재시작 세션 꼬임 완화).""" with self._reg_timer_lock: if self._reg_timer: try: self._reg_timer.cancel() except Exception: pass self._reg_timer = None try: self._unsubscribe_all_before_close(reason="stop", clear_ram=True) except Exception as e: logger.warning("키움 WS 종료 전 REMOVE 실패: %s", e) self._running = False self._connected = False self._authenticated = False try: if self._ws is not None: self._ws.close() except Exception: pass logger.info("⏹ 키움 WS 종료") def _max_subscriptions(self) -> int: """그룹당 최대 구독 수 — DB ``KIWOOM_WS_MAX_SUBSCRIPTIONS`` (기본 100).""" v = get_env_int("KIWOOM_WS_MAX_SUBSCRIPTIONS", self.MAX_SUBSCRIPTIONS_PER_GROUP) return max(1, min(200, int(v))) def _reg_chunk_size(self) -> int: """REG 한 메시지당 최대 종목 수 — ``KIWOOM_WS_REG_CHUNK_SIZE`` (기본 25).""" v = get_env_int("KIWOOM_WS_REG_CHUNK_SIZE", 25) return max(1, min(80, int(v))) def _reg_gap_sec(self) -> float: """청크 사이 전송 간격(초) — ``KIWOOM_WS_REG_GAP_SEC`` (기본 0.18).""" return max(0.05, min(2.0, get_env_float("KIWOOM_WS_REG_GAP_SEC", 0.18))) def _reg_debounce_sec(self) -> float: """단건 subscribe 묶기 대기(초) — ``KIWOOM_WS_REG_DEBOUNCE_SEC`` (기본 0.12).""" return max(0.0, min(1.0, get_env_float("KIWOOM_WS_REG_DEBOUNCE_SEC", 0.12))) def _orderbook_ws_enabled(self) -> bool: # 사용자의 혼동 방지를 위해 KIWOOM_WS_ORDERBOOK_ENABLED 대신, # LIVE_OB_PROVIDER가 kiwoom이거나, 키움 호가 적재(WS_ORDERBOOK_SAVE_KIWOOM)가 켜져있으면 자동 수신. from ..utils.env import get_env_from_db, get_env_bool live_ob_provider = (get_env_from_db("LIVE_OB_PROVIDER", "kiwoom") or "kiwoom").strip().lower() save_kiwoom = get_env_bool("WS_ORDERBOOK_SAVE_KIWOOM", True) if live_ob_provider == "kiwoom" or save_kiwoom: return True return get_env_bool("KIWOOM_WS_ORDERBOOK_ENABLED", False) def _program_ws_enabled(self) -> bool: return get_env_bool("KIWOOM_WS_PROGRAM_ENABLED", True) def is_available(self) -> bool: """websocket-client 설치 및 키 설정 여부.""" return bool(self._available and self.app_key and self.app_secret) def is_authenticated(self) -> bool: """LOGIN OK 이후 REG/조건검색 전송 가능.""" return bool(self._connected and self._authenticated) def register_trnm_handler(self, trnm: str, fn: Callable) -> None: """외부 모듈(조건검색 CNSR* 등)용 trnm 핸들러 — 시세 WS 와 세션 공유.""" key = str(trnm or "").strip().upper() if not key or not callable(fn): return with self._ext_handler_lock: if fn not in self._ext_handlers[key]: self._ext_handlers[key].append(fn) def unregister_trnm_handler(self, trnm: str, fn: Callable) -> None: key = str(trnm or "").strip().upper() with self._ext_handler_lock: lst = self._ext_handlers.get(key) if lst and fn in lst: lst.remove(fn) def add_on_login_callback(self, fn: Callable) -> None: """LOGIN OK 직후(재접속마다) 호출 — 조건검색 CNSRLST 등.""" if not callable(fn): return with self._login_cb_lock: if fn not in self._login_callbacks: self._login_callbacks.append(fn) def remove_on_login_callback(self, fn: Callable) -> None: with self._login_cb_lock: if fn in self._login_callbacks: self._login_callbacks.remove(fn) def send_json(self, msg: dict) -> bool: """인증된 WS 에 JSON 전송 (조건검색 CNSRREQ 등).""" if not self.is_authenticated() or not self._ws: log_kiwoom_ws( logger, "send_skip", level="warning", trnm=str(msg.get("trnm") or ""), extra="미인증 또는 ws=None", force=True, ) return False try: self._ws.send(json.dumps(msg)) return True except Exception as e: log_kiwoom_ws( logger, "send_fail", level="warning", trnm=str(msg.get("trnm") or ""), extra=str(e), force=True, ) logger.debug("키움 WS send_json 실패: %s", e) return False def _fire_login_callbacks(self, ws) -> None: with self._login_cb_lock: cbs = list(self._login_callbacks) for cb in cbs: try: cb(ws) except Exception as e: logger.debug("키움 WS login callback 예외: %s", e) def _dispatch_ext_handlers(self, trnm: str, ws, msg: dict) -> None: key = str(trnm or "").strip().upper() with self._ext_handler_lock: handlers = list(self._ext_handlers.get(key, [])) for fn in handlers: try: fn(ws, msg) except Exception as e: logger.debug("키움 WS ext handler(%s) 예외: %s", key, e) def _reg_types(self) -> List[str]: """REG/REMOVE 실시간 타입 — 0B 체결 + (옵션) 0D 호가 + 0w 프로그램.""" types = [self.SUB_TYPE] if self._orderbook_ws_enabled(): types.append(self.SUB_TYPE_ORDERBOOK) if self._program_ws_enabled(): types.append(self.SUB_TYPE_PROGRAM) return types def subscribe_many(self, codes: Iterable[str]) -> List[str]: """여러 종목 등록. REG 는 청크+간격으로 전송(TRNM 레이트리밋 회피). LOGIN 전이면 집합만 갱신하고, LOGIN OK 시 ``_send_reg_chunked`` 로 일괄 전송. """ added: List[str] = [] with self._sub_lock: for raw in codes: c = (str(raw) or "").strip() if not c or c in self._subscribed: continue if len(self._subscribed) >= self._max_subscriptions(): logger.warning( "⚠️ 키움 WS 구독 한도 초과 (%d/%d) — 이후 종목 스킵", len(self._subscribed), self._max_subscriptions(), ) break self._subscribed.add(c) added.append(c) if added and self._connected and self._authenticated: self._send_reg_chunked(added) return added def subscribe(self, code: str) -> bool: """단일 종목 등록. REG 는 짧게 디바운스 후 묶어 전송.""" 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(): logger.warning( "⚠️ 키움 WS 구독 한도 초과 (%d/%d) — %s 등록 거절", len(self._subscribed), self._max_subscriptions(), code, ) return False self._subscribed.add(code) # LOGIN 전에는 집합만 쌓고, REG 는 LOGIN OK 후 일괄 전송(중복·순서 레이스 방지) need_reg = self._connected and self._authenticated if need_reg: self._reg_batch_codes.add(code) if need_reg: self._schedule_reg_debounce() return True def _schedule_reg_debounce(self) -> None: deb = self._reg_debounce_sec() if deb <= 0: with self._sub_lock: pending = sorted(self._reg_batch_codes) self._reg_batch_codes.clear() if pending: self._send_reg_chunked(pending) return with self._reg_timer_lock: if self._reg_timer: try: self._reg_timer.cancel() except Exception: pass self._reg_timer = threading.Timer(deb, self._flush_reg_debounced) self._reg_timer.daemon = True self._reg_timer.start() def _flush_reg_debounced(self) -> None: with self._reg_timer_lock: self._reg_timer = None with self._sub_lock: pending = sorted(self._reg_batch_codes) self._reg_batch_codes.clear() if pending and self._connected and self._authenticated: self._send_reg_chunked(pending) 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._orderbook_cache: self._orderbook_cache.remove(code) if self._program_cache: self._program_cache.remove(code) if self._connected and self._authenticated: return self._send_remove([code]) return True def get_orderbook_snapshot( self, code: str, max_age_sec: float = 3.0, ) -> Optional["OrderbookSnapshot"]: """키움 0D RAM 호가 스냅샷.""" if not self._orderbook_cache: return None return self._orderbook_cache.get(code, max_age_sec=max_age_sec) def get_orderbook(self, code: str, max_age_sec: float = 3.0) -> Optional[Dict]: """KIS REST 호가 dict 호환 — ``orderbook_sell`` / OrderManager 폴백용.""" if not self._orderbook_cache: return None return self._orderbook_cache.get_kis_bid_dict(code, max_age_sec=max_age_sec) def get_program_snapshot( self, code: str, max_age_sec: float = 30.0, ) -> Optional["ProgramSnapshot"]: """키움 0w RAM 프로그램매매 스냅샷.""" if not self._program_cache: return None return self._program_cache.get(code, max_age_sec=max_age_sec) def add_price_listener(self, callback) -> None: """현재가 갱신 콜백 등록. callback(code, price, data_dict).""" if callback is None: return with self._price_listener_lock: if callback not in self._price_listeners: self._price_listeners.append(callback) def remove_price_listener(self, callback) -> None: if callback is None: return with self._price_listener_lock: try: self._price_listeners.remove(callback) except ValueError: pass def _emit_price_listeners(self, code: str, price: float, data: Dict) -> None: with self._price_listener_lock: listeners = list(self._price_listeners) for cb in listeners: try: cb(code, price, data) except Exception as ex: logger.debug("price_listener 예외 %s: %s", code, ex) def get_price(self, code: str, max_age_sec: Optional[float] = None) -> Optional[Dict]: """KIS WS ``get_price`` 와 동일 포맷 반환. None = 마지막 RAM (체결 공백이어도 유지). 반환:: { "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 max_age_sec is not None and 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 attach_candle_aggregator(self, agg: Any) -> None: """KIS ``CandleAggregator`` 연결 — 키움 0B 틱으로 분봉 RAM 집계. ``set_candle_tick_codes()`` 로 후보 종목만 필터링하지 않으면 구독된 모든 종목 틱이 집계기로 들어감 (검증 전용 모드에서는 필터 권장). """ self._candle_agg = agg logger.info("✅ KiwoomWebSocket → CandleAggregator 연결") def set_candle_tick_codes(self, codes: Optional[Set[str]]) -> None: """집계기로 보낼 종목 코드. None=전체 구독 종목, 비어있지 않은 Set=해당 코드만.""" self._candle_agg_codes = codes def attach_tick_recorder(self, recorder: Any) -> None: """TickRecorder 연결 — 키움 0B 체결 틱.""" self._tick_recorder = recorder logger.info("✅ KiwoomWebSocket → TickRecorder 연결") def attach_trigger_snapshot_recorder(self, recorder: Any) -> None: """TriggerSnapshotRecorder 연결 — 키움 0D/0w 스냅샷 DB.""" self._trigger_snapshot_recorder = recorder logger.info("✅ KiwoomWebSocket → TriggerSnapshotRecorder 연결") # ------------------------------------------------------------------ # 내부: 메인 수신 루프 # ------------------------------------------------------------------ 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)") try: from kis_trader.utils.ops_alert import ops_alert ops_alert( "token_kiwoom", "키움 토큰 발급 실패", detail="시세 WS·조건검색 불가 → 60s 후 재시도", level="critical", session_only=False, ) except Exception: pass 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 — 연결 종료까지 여기서 대기 # ping_interval=0 : 키움은 JSON {"trnm":"PING"} keep-alive (프로토콜 ping 비호환) self._ws.run_forever(ping_interval=0) 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 in ("CNSRLST", "CNSRREQ", "CNSRCLR", "REG", "REMOVE", "LOGIN"): log_kiwoom_api_msg(logger, str(trnm or "?"), msg) 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 # 안정 LOGIN 성공 → 재연결 카운터 리셋 (주말 Bye 후 소모분 복구) self._reconnect_count = 0 logger.info("✅ 키움 WS LOGIN OK") # 디바운스 타이머 취소 — LOGIN 직후 일괄 REG 와 이중 전송 방지 with self._reg_timer_lock: if self._reg_timer: try: self._reg_timer.cancel() except Exception: pass self._reg_timer = None # 누적된 구독 일괄 등록 (청크+간격 — TRNM=REG 레이트리밋) with self._sub_lock: pending = sorted(self._subscribed) self._reg_batch_codes.clear() if pending: # LOGIN ack 를 받은 같은 on_message 스레드에서 REG 청크 전송. # 이 동안 0B 파싱이 밀림. 장중 추가는 Timer 디바운스(별 스레드). self._send_reg_chunked(pending, group_replace=True) # 조건검색 등 공유 세션 모듈 — LOGIN 직후 CNSRLST 재등록 self._fire_login_callbacks(ws) else: logger.warning("❌ 키움 WS LOGIN 실패 rc=%s msg=%s", rc, rm) try: from kis_trader.utils.ops_alert import ops_alert ops_alert( "ws_kiwoom_login", "키움 WS LOGIN 실패", detail=f"rc={rc} msg={rm}", level="critical", session_only=False, ) except Exception: pass # 805004 등 서버 토큰 거부 → 캐시 만료 전이라도 무효화 후 재발급 if self._is_token_rejected_login(rc, rm): self._invalidate_token_cache( "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 실패 rc=%s msg=%s", rc, msg.get("return_msg", ""), ) return if trnm == "REMOVE": return # 무관심 if trnm == "REAL": self._handle_real(msg) self._dispatch_ext_handlers("REAL", ws, msg) return # 조건검색 CNSRLST / CNSRREQ / CNSRCLR 등 self._dispatch_ext_handlers(trnm, ws, msg) def _handle_real(self, msg: dict) -> None: """실시간 데이터 처리 — 0B 체결 + 0D 호가잔량 + 0w 프로그램매매.""" items = msg.get("data") or [] for item in items: sub_type = str(item.get("type", "")).strip() code = str(item.get("item", "")).strip() values = item.get("values") or {} if sub_type == self.SUB_TYPE: self._cache_tick(code, values) elif sub_type == self.SUB_TYPE_ORDERBOOK and self._orderbook_cache: try: snap = self._orderbook_cache.update_from_kiwoom_0d(code, values) if snap is not None and self._trigger_snapshot_recorder is not None: tt = str(values.get("20") or values.get(self.FID_TICK_TIME) or "").strip() self._trigger_snapshot_recorder.on_orderbook(snap, snap_time=tt or None) except Exception as ex: logger.debug("키움 0D 파싱 실패 %s: %s", code, ex) elif sub_type == self.SUB_TYPE_PROGRAM and self._program_cache: try: snap = self._program_cache.update_from_kiwoom_0w(code, values) if self._trigger_snapshot_recorder is not None: tt = str(values.get("20") or "").strip() self._trigger_snapshot_recorder.on_program(snap, snap_time=tt or None) except Exception as ex: logger.debug("키움 0w 파싱 실패 %s: %s", code, ex) 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 호환 + 키움 0B 원본 FID 보존 (분석·추후 필터용) exec_raw = values.get(self.FID_EXEC_STRENGTH, "0") try: cntr_str = float(str(exec_raw).replace(",", "")) except (ValueError, TypeError): cntr_str = 0.0 upper_limit_raw = str(values.get(self.FID_UPPER_LIMIT_TIME, "") or "").strip() lower_limit_raw = str(values.get(self.FID_LOWER_LIMIT_TIME, "") or "").strip() 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")), "cntr_str": str(cntr_str), "upper_limit_time": upper_limit_raw, "lower_limit_time": lower_limit_raw, "kiwoom_fid228": str(exec_raw), "kiwoom_fid567": upper_limit_raw, "kiwoom_fid568": lower_limit_raw, "chetime": "", "kiwoom_fid20": "", } # tick_time / tick_vol — CandleAggregator·TickRecorder 공용 (보유-only도 recorder 수집) # 늦음 판정 = 이 틱의 체결시각(FID20) vs 우리 서버 지금. 증권사끼리 비교 아님. # 매매 RAM/리스너/봉: 읽기나이(LIVE_FEED_FALLBACK, 기본 3초) 넘으면 안 넣음 → 2차·3차가 연다. # 적재(ws_ticks): 버리지 않고 전부. LIVE_MAX<=0 이면 적재 skip 없음. # 아침 FID20 동결: 같은 FID20 가 TIME_MAX초 이상 안 움직일 때만 RAM 유지. import datetime as _dt tt_raw = str(values.get(self.FID_TICK_TIME, "") or "").strip() _now_dt = _dt.datetime.now() _now_wall = time.time() if len(tt_raw) >= 6: tick_time_pkt = tt_raw[-6:] else: tick_time_pkt = _now_dt.strftime("%H%M%S") data_compat["chetime"] = tick_time_pkt data_compat["kiwoom_fid20"] = tick_time_pkt data_compat["tick_time"] = ( tt_raw[:14] if len(tt_raw) >= 14 else (_now_dt.strftime("%Y%m%d") + tick_time_pkt) ) _pkt_dt = None try: if len(tt_raw) >= 14: _pkt_dt = _dt.datetime.strptime(tt_raw[:14], "%Y%m%d%H%M%S") elif len(tt_raw) >= 6: _pkt_dt = _dt.datetime.strptime( _now_dt.strftime("%Y%m%d") + tick_time_pkt, "%Y%m%d%H%M%S", ) except Exception: _pkt_dt = None _lag_sec = 0.0 if _pkt_dt is not None: _lag_sec = float((_now_dt - _pkt_dt).total_seconds()) try: _live_max = int(get_env_int("KIWOOM_TICK_LIVE_MAX_LAG_SEC", 0) or 0) except (TypeError, ValueError): _live_max = 0 try: from kis_trader.engine.feed_fallback import live_feed_fallback_max_age_sec _read_max = float(live_feed_fallback_max_age_sec()) except Exception: _read_max = 2.0 _time_max = max(5, int(get_env_int("KIWOOM_TICK_TIME_MAX_LAG_SEC", 120))) _last_fid = self._fid20_last.get(code) if _pkt_dt is not None and _last_fid != _pkt_dt: self._fid20_first_wall[code] = _now_wall _first_w = float(self._fid20_first_wall.get(code, _now_wall) or _now_wall) # 같은 FID20 가 TIME_MAX초 이상 지속 = 아침 시계 멈춤. 밀린 줄의 같은 초 2건(수십 ms)은 제외. _frozen = bool( _pkt_dt is not None and _last_fid is not None and _pkt_dt == _last_fid and (_now_wall - _first_w) >= float(_time_max) ) # 재연결 버퍼: FID20가 지금보다 읽기나이(2초) 이상 과거면 매매 RAM 미반영. # 적재는 기본 전부(LIVE_MAX=0). 양수일 때만 그 초 초과 미저장. _skip_ram = bool( _read_max > 0 and _pkt_dt is not None and _lag_sec > float(_read_max) and not _frozen ) _skip_persist = bool( _live_max > 0 and _pkt_dt is not None and _lag_sec > float(_live_max) and not _frozen ) if _pkt_dt is not None: self._fid20_last[code] = _pkt_dt if not _skip_ram: wall_ts = time.time() with self._cache_lock: self._cache[code] = {"data": data_compat, "ts": wall_ts} self._emit_price_listeners(code, float(price), data_compat) else: _log_gap = max(5.0, float(get_env_float("KIWOOM_TICK_STALE_LOG_GAP_SEC", 30.0) or 30.0)) _prev_log = float(self._stale_skip_log_ts.get(code, 0.0) or 0.0) if (time.time() - _prev_log) >= _log_gap: self._stale_skip_log_ts[code] = time.time() logger.warning( "키움 0B stale skip RAM %s fid20=%s lag=%.0fs price=%.0f (매수체크 미반영·2차 폴백)", code, tick_time_pkt, _lag_sec, price, ) tick_time = tick_time_pkt if _pkt_dt is None: tick_time = _now_dt.strftime("%H%M%S") elif _frozen and _lag_sec > float(_time_max): # FID20 동결(아침 55분 버그) — 분봉·틱청산 타임라인만 wall-clock tick_time = _now_dt.strftime("%H%M%S") try: tick_vol = int( abs(float(str(values.get(self.FID_TICK_VOL, "0")).replace(",", ""))) ) except (ValueError, TypeError): tick_vol = 0 # ── CandleAggregator: 후보만 (tick_to_agg = 후보−보유) ───── # 밀린 메인 틱은 확정봉에 넣지 않음 (freeze 오염 방지) if self._candle_agg is not None and not _skip_ram: filt = self._candle_agg_codes if filt is None or code in filt: try: self._candle_agg.on_tick(code, price, tick_vol, tick_time, source="kiwoom") except Exception as ex: logger.debug("키움→CandleAggregator on_tick 실패 %s: %s", code, ex) # ── TickRecorder: 보유 포함 전 구독 종목 (백테 틱청산·B안 폴백) ───── # 공책은 밀려도 적재. 읽기 2초와 분리. 옵투나는 lag>2초 메인을 보조로 교체. if self._tick_recorder is not None and not _skip_persist: try: self._tick_recorder.on_tick( code, price, tick_vol, tick_time, source="kiwoom", cntr_str=cntr_str, upper_limit_time=upper_limit_raw or None, tick_time_raw=tt_raw or None, ) except Exception as ex: logger.debug("키움→TickRecorder on_tick 실패 %s: %s", code, ex) # ── 호가 틱동기: 체결 1건당 RAM 0D 스냅 1장 (LS tick 모드와 동일) ── # 밀린 체결시각에 지금 호가를 붙이면 옵투나 호가가 미래가 됨 → RAM과 같이 skip. if self._trigger_snapshot_recorder is not None and not _skip_ram: try: max_age = float( getattr( self._trigger_snapshot_recorder, "orderbook_tick_max_age_sec", 3.0, ) or 3.0 ) snap = self.get_orderbook_snapshot(code, max_age_sec=max_age) if snap is not None: self._trigger_snapshot_recorder.on_orderbook_tick_sync( code, snap, snap_time=tick_time, ) except Exception as ex: logger.debug("키움→호가틱동기 실패 %s: %s", code, ex) 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 err_s = str(error or "") log_kiwoom_ws( logger, "on_error", level="warning", extra=err_s, force=True, ) logger.warning("⚠️ 키움 WS 오류: %s", error) def _on_close(self, ws, close_status_code, close_msg) -> None: self._connected = False self._authenticated = False log_kiwoom_ws( logger, "on_close", level="info", rc=close_status_code, msg=close_msg or "", force=True, ) logger.info("🔌 키움 WS 연결 종료 (code=%s msg=%s)", close_status_code, close_msg or "") # ------------------------------------------------------------------ # 내부: 등록/해지 메시지 발송 # ------------------------------------------------------------------ def _reg_refresh_keep(self) -> str: """REG refresh. 기본 1(기존 유지). 0=그룹 기존 item/type 해지 — 테스트용.""" v = int(get_env_int("KIWOOM_WS_REG_REFRESH", 1) or 1) return "0" if v == 0 else "1" def _send_reg_chunked(self, codes: List[str], *, group_replace: bool = False) -> None: """REG 를 청크 단위로 전송하고 청크 사이에 sleep (키움 TRNM 건수 제한 회피). group_replace=True (LOGIN 일괄): ``KIWOOM_WS_REG_REFRESH=0`` 이면 **첫 청크만** refresh=0. 이후 청크·장중 추가 REG는 항상 1 (0을 반복하면 앞 청크 구독이 삭제됨). """ if not codes: return chunk = self._reg_chunk_size() gap = self._reg_gap_sec() want = self._reg_refresh_keep() if group_replace else "1" if group_replace and want == "0": logger.warning( "⚠️ 키움 WS REG refresh=0 테스트 — LOGIN 일괄 첫 청크만 기존 구독 해지 후 재등록" ) for i in range(0, len(codes), chunk): part = codes[i:i + chunk] refresh = want if i == 0 else "1" self._send_reg(part, refresh=refresh) if i + chunk < len(codes): time.sleep(gap) def _send_reg(self, codes: list, *, refresh: str = "1") -> bool: """REG 발송 (그룹 1, 0B+0D 타입). 한 번에 여러 종목 OK.""" if not codes or not self._ws or not self._connected: return False flag = "0" if str(refresh).strip() == "0" else "1" try: self._ws.send(json.dumps({ "trnm": "REG", "grp_no": self.GROUP_NO, "refresh": flag, "data": [{"item": list(codes), "type": self._reg_types()}], })) logger.info( "📡 키움 WS REG 발송: %d종목 types=%s refresh=%s (총 %d/%d)", len(codes), self._reg_types(), flag, self.subscribed_count(), self._max_subscriptions(), ) 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 or not self._connected: return False try: self._ws.send(json.dumps({ "trnm": "REMOVE", "grp_no": self.GROUP_NO, "data": [{"item": list(codes), "type": self._reg_types()}], })) 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 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) @staticmethod def _is_token_rejected_login(rc: Any, rm: Any) -> bool: """WS LOGIN 응답이 '토큰 무효'인지 — 캐시 invalidate 대상.""" try: if int(rc) == 805004: return True except Exception: pass m = str(rm or "") return ("Token이 유효하지 않습니다" in m) or ("토큰이 유효하지 않습니다" in m) def _invalidate_token_cache(self, reason: str) -> None: try: from .kis_ws import _invalidate_kiwoom_token_cached except ImportError: logger.warning("kis_ws._invalidate_kiwoom_token_cached import 실패") return _invalidate_kiwoom_token_cached( self.app_key, self.app_secret, self.is_mock, reason=reason, ) # ------------------------------------------------------------------ # 내부: 재연결 백오프 (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 # 시간당 한도 초과 → 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]) log_kiwoom_ws( logger, "reconnect_hourly_cap", level="warning", extra=f"wait={int(wait)}s count={len(self._reconnect_times)}", force=True, ) logger.warning("⚠️ 키움 WS 시간당 재연결 한도 초과 → %ds 대기", int(wait)) time.sleep(max(60.0, wait)) return # 총 한도 초과 — 기본은 긴 쿨다운 후 카운터 리셋 (주말 토큰거부 연타로 # 월요일까지 영구 OFF 되면 SCALP/MOMENTUM 조건식 REAL 이 끊김). # 응급으로만 KIWOOM_WS_HARD_DISABLE_ON_MAX_RECONNECT=true. if self._reconnect_count >= self.MAX_RECONNECT_ATTEMPTS: if get_env_bool("KIWOOM_WS_HARD_DISABLE_ON_MAX_RECONNECT", False): logger.warning("⛔ 키움 WS 총 재연결 한도 초과 → 자동 비활성") self._running = False return cool = float(get_env_float("KIWOOM_WS_MAX_RECONNECT_COOLDOWN_SEC", 1800.0)) cool = max(60.0, min(cool, 86400.0)) logger.warning( "⚠️ 키움 WS 총 재연결 한도(#%d) → %.0fs 쿨다운 후 카운터 리셋 " "(영구 OFF 안 함 · HARD_DISABLE=true 시만 비활성)", self._reconnect_count, cool, ) # 쿨다운 전 토큰 캐시 비우기 — 죽은 토큰 재전송 루프 차단 self._invalidate_token_cache("max reconnect cooldown") time.sleep(cool) self._reconnect_count = 0 self._reconnect_times.clear() return attempt = self._reconnect_count + 1 delay = ws_reconnect_delay_for_attempt(attempt) logger.info("⏳ 키움 WS %ds 후 재연결 시도 (#%d)", int(delay), attempt) time.sleep(delay) self._reconnect_times.append(time.time()) self._reconnect_count += 1