""" kis_trader/network/kiwoom_condition_manager.py — 키움 조건검색 기반 동적 유니버스 (WS 실시간) ================================================================================================== 팩트 체크 먼저: * KIS 조건검색(psearch-title/psearch-result)은 **웹소켓 미지원 → REST 폴링만** 가능하다. (기존 ``ConditionSearchManager`` 참고) * 반면 키움 신형 오픈API 는 **웹소켓으로 조건검색 실시간(편입/이탈 push)** 을 지원한다. - CNSRLST : 서버 저장 조건식 목록 (name → seq 해결) - CNSRREQ : 실시간 조건검색 등록 (search_type="1") → 초기 매칭 + 이후 REAL push - REAL : 편입(843="I") / 이탈(843="D"), 종목코드는 9001 - CNSRCLR : 실시간 해제 (이 경로는 ``_test_kiwoom_condition_realtime.py`` 에서 실계정으로 검증 완료) 설계 (최소 침습): * ``ConditionSearchManager`` 를 **서브클래싱** 하여 - 결과 반영 로직(``_apply_result``), 순서 계산(``_build_ordered_universe``), 스냅샷 저장(``_save_snapshot``), 조회 API(``get_universe_for``/``get_candidates_for``), EXIT grace, name_map, _configs 정규화 등을 **그대로 재사용**한다. - 데이터 취득만 REST 폴링 → **키움 WS 실시간** 으로 오버라이드. * 스냅샷은 KIS 와 **동일한** ``target_candidates_history`` 에 저장하므로 백테스트/유니버스 타임라인 코드는 변경 없이 그대로 재현 가능하다. * BaseStrategy 는 ``kiwoom_condition_mgr`` 로 주입받아 ``{SID}_UNIVERSE_SOURCE=kiwoom_condition`` 일 때 이 매니저를 소비한다. (KIS ``condition`` 소스와 완전 독립 — 기존 동작 불변) 사용: km = KiwoomConditionSearchManager( app_key=..., app_secret=..., # 반드시 KIWOOM_APP_KEY_REAL (main 이 is_mock=False 고정) configs=[{"strategy_id": "MOMENTUM", "name": "momentum", "seq": "3"}], db=db, ) km.start() codes = km.get_universe_for("MOMENTUM") """ from __future__ import annotations import json import random import threading import time from typing import Any, Dict, List, Optional, Set from .condition_manager import ConditionSearchManager from ..utils.env import get_env_float, get_env_from_db, get_env_int from ..utils.logger import get_logger from ..ws.kis_ws import _get_kiwoom_token_cached logger = get_logger("kis_trader.kwcond") def _normalize_code(raw) -> str: """키움 종목코드 정규화: 'A005930' → '005930' (KIS 6자리 코드계와 정합).""" c = str(raw or "").strip() if not c: return "" # 키움은 국내주식 코드 앞에 'A' 접두를 붙이는 경우가 있음. if c[0] in ("A", "a") and len(c) >= 7: c = c[1:] return c class KiwoomConditionSearchManager(ConditionSearchManager): """키움 웹소켓 실시간 조건검색 매니저. ``ConditionSearchManager`` 와 **동일한 public API** 를 제공한다: - start() / stop() - get_universe_for(strategy_id) / get_candidates_for(strategy_id) - _configs (BaseStrategy._is_strategy_registered 판별용) """ def __init__( self, *, app_key: str, app_secret: str, is_mock: bool, configs: Optional[List[Dict]] = None, db=None, on_change=None, shared_ws: Any = None, ): # 부모 초기화: client 는 REST 미사용이므로 None, user_id 는 로깅용 placeholder. # configs 정규화·EXIT grace·name_map·_lock 등은 부모가 세팅. super().__init__( client=None, user_id="KIWOOM", configs=configs, db=db, on_change=on_change, ) self._app_key = (app_key or "").strip() self._app_secret = (app_secret or "").strip() self._is_mock = bool(is_mock) self._token: Optional[str] = None # 시세 WS(KiwoomWebSocketPriceCache) 와 **단일 세션 공유** — 별도 접속 시 Bye 루프 self._shared_ws: Any = shared_ws self._shared_mode: bool = False self._shared_handlers_bound: bool = False # WS URL (실전/모의) — env 로 오버라이드 가능. if self._is_mock: self._ws_url = ( get_env_from_db( "KIWOOM_WS_URL_MOCK", "wss://mockapi.kiwoom.com:10000/api/dostk/websocket", ) or "wss://mockapi.kiwoom.com:10000/api/dostk/websocket" ).strip() else: self._ws_url = ( get_env_from_db( "KIWOOM_WS_URL_REAL", "wss://api.kiwoom.com:10000/api/dostk/websocket", ) or "wss://api.kiwoom.com:10000/api/dostk/websocket" ).strip() # 실시간 재등록 대기(초) — 재접속 시 사용. self._reconnect_backoff = float(get_env_int("KIWOOM_COND_RECONNECT_SEC", 5)) # 키움 전용 상태 self._kw_lock = threading.Lock() self._ws = None self._ws_thread: Optional[threading.Thread] = None # seq → {code(정규화): name} (삽입순 유지 = HTS 응답 순서) self._seq_codes: Dict[str, "Dict[str, str]"] = {} # seq → [strategy_id, ...] (같은 seq 를 공유하는 전략) self._sid_by_seq: Dict[str, List[str]] = {} # 이번 접속에서 실시간 등록할 unique seq 목록 (CNSRLST 해결 후 채움) self._active_seqs: List[str] = [] # 최초 CNSRLST 처리 + CNSRREQ 시도 완료 신호 (start() 동기 대기용) self._ready = threading.Event() self._start_ok = False # CNSRREQ 발송·응답 추적 (연속 발송 시 응답 누락 → 재발송) self._cnsrreq_pending: Set[str] = set() self._cnsrreq_confirmed: Set[str] = set() self._cnsrreq_retry_timer: Optional[threading.Timer] = None # ------------------------------------------------------------------ # Public API (오버라이드) — 부모 start() 는 REST 폴링이므로 사용 안 함 # ------------------------------------------------------------------ def start(self) -> bool: """토큰 발급 → 키움 WS 접속 → CNSRLST 로 seq 해결 → CNSRREQ 실시간 등록.""" if not self._configs: logger.info("키움 조건검색 configs 비어 있음 → 매니저 비활성") return False if not (self._app_key and self._app_secret): logger.warning("키움 앱키/시크릿 누락 → 조건검색 매니저 비활성") return False # 시세 WS 가 이미 떠 있으면 **같은 소켓**으로 조건검색 (키움 1세션 정책) if self._shared_ws is not None and getattr(self._shared_ws, "is_available", lambda: False)(): return self._start_shared() try: import websocket # noqa: F401 (websocket-client 존재 확인) except Exception as e: logger.warning("websocket-client 미설치 → 키움 조건검색 비활성: %s", e) return False return self._start_own_connection() def _start_shared(self) -> bool: """KiwoomWebSocketPriceCache 세션에 CNSR* 핸들러만 부착 (별도 접속 없음).""" self._shared_mode = True self._running = True self._token = _get_kiwoom_token_cached( self._app_key, self._app_secret, self._is_mock ) if not self._token: logger.warning("키움 토큰 발급 실패 → 조건검색 매니저 비활성") return False self._bind_shared_handlers() # 이미 LOGIN 된 상태면 즉시 CNSRLST if getattr(self._shared_ws, "is_authenticated", lambda: False)(): self._send_cnsrlst() ready_timeout = float(get_env_int("KIWOOM_COND_START_TIMEOUT_SEC", 10)) self._ready.wait(timeout=ready_timeout) if self._start_ok: logger.info( "✅ 키움 조건검색 실시간 시작 [공유WS] (%d개, mock=%s, exit_grace=%ds, history=%s)", len(self._active_seqs), self._is_mock, int(self._exit_grace_sec), "ON" if (self.history_enabled and self.db is not None) else "OFF", ) else: logger.warning( "⚠️ 키움 조건검색 [공유WS] 초기 등록 미완료(타임아웃) — LOGIN 재접속 시 자동 재시도" ) return True def _start_own_connection(self) -> bool: """레거시: 단독 WS (shared_ws 없을 때만 — 중복 접속 주의).""" self._token = _get_kiwoom_token_cached( self._app_key, self._app_secret, self._is_mock ) if not self._token: logger.warning("키움 토큰 발급 실패 → 조건검색 매니저 비활성") return False self._running = True self._ws_thread = threading.Thread( target=self._run_ws_forever, daemon=True, name="KwCondSearch" ) self._ws_thread.start() ready_timeout = float(get_env_int("KIWOOM_COND_START_TIMEOUT_SEC", 10)) self._ready.wait(timeout=ready_timeout) if self._start_ok: logger.info( "✅ 키움 조건검색 실시간 시작 (%d개, mock=%s, exit_grace=%ds, history=%s)", len(self._active_seqs), self._is_mock, int(self._exit_grace_sec), "ON" if (self.history_enabled and self.db is not None) else "OFF", ) else: logger.warning( "⚠️ 키움 조건검색 실시간 초기 등록 미완료(타임아웃) — 백그라운드 재시도 지속" ) return True def stop(self) -> None: self._running = False if self._shared_mode: self._unbind_shared_handlers() return try: if self._ws is not None: self._ws.close() except Exception: pass def _bind_shared_handlers(self) -> None: if self._shared_handlers_bound or not self._shared_ws: return ws = self._shared_ws ws.register_trnm_handler("CNSRLST", self._on_shared_trnm) ws.register_trnm_handler("CNSRREQ", self._on_shared_trnm) ws.register_trnm_handler("CNSRCLR", self._on_shared_trnm) ws.register_trnm_handler("REAL", self._on_shared_real) ws.add_on_login_callback(self._on_shared_login) self._shared_handlers_bound = True logger.info("🔗 키움 조건검색 → 시세 WS 세션 공유 (중복 접속 방지)") def _unbind_shared_handlers(self) -> None: if not self._shared_handlers_bound or not self._shared_ws: return ws = self._shared_ws ws.unregister_trnm_handler("CNSRLST", self._on_shared_trnm) ws.unregister_trnm_handler("CNSRREQ", self._on_shared_trnm) ws.unregister_trnm_handler("CNSRCLR", self._on_shared_trnm) ws.unregister_trnm_handler("REAL", self._on_shared_real) ws.remove_on_login_callback(self._on_shared_login) self._shared_handlers_bound = False def _send_cnsrlst(self) -> None: if self._shared_ws: self._shared_ws.send_json({"trnm": "CNSRLST"}) def _on_shared_login(self, ws) -> None: """시세 WS 재접속마다 조건식 목록 재조회 → CNSRREQ 재등록.""" if not self._running: return self._ready.clear() self._start_ok = False try: ws.send(json.dumps({"trnm": "CNSRLST"})) logger.debug("키움 조건검색 [공유WS] LOGIN → CNSRLST") except Exception as e: logger.debug("키움 조건검색 CNSRLST 발송 실패: %s", e) def _on_shared_trnm(self, ws, msg: dict) -> None: trnm = msg.get("trnm") if trnm == "CNSRLST": self._handle_condition_list(ws, msg.get("data") or []) elif trnm == "CNSRREQ": self._handle_cnsrreq(msg) elif trnm == "CNSRCLR": logger.debug("키움 조건검색 CNSRCLR 응답: rc=%s", msg.get("return_code")) def _on_shared_real(self, ws, msg: dict) -> None: """조건검색 편입/이탈 REAL — 843 필드 있는 항목만 처리.""" rows = msg.get("data") or [] cond_rows = [] for it in rows: vals = it.get("values") if isinstance(it, dict) else None if isinstance(vals, dict) and "843" in vals: cond_rows.append(it) if cond_rows: self._handle_real(cond_rows) # ------------------------------------------------------------------ # WS 라이프사이클 # ------------------------------------------------------------------ def _run_ws_forever(self) -> None: """접속 → (끊기면) 백오프 후 재접속 루프. 재접속 시 실시간 재등록.""" import websocket # websocket-client while self._running: try: self._ws = websocket.WebSocketApp( self._ws_url, on_open=self._on_open, on_message=self._on_message, on_error=self._on_error, on_close=self._on_close, ) # ping_interval=0 : 키움은 서버가 PING 을 보내면 echo 하는 방식. self._ws.run_forever(ping_interval=0) except Exception as e: logger.debug("키움 조건검색 WS 예외: %s", e) if not self._running: break # 재접속 대기 (중단 감지 해상도 0.5s) deadline = time.time() + self._reconnect_backoff while self._running and time.time() < deadline: time.sleep(0.5) def _on_open(self, ws) -> None: try: ws.send(json.dumps({"trnm": "LOGIN", "token": self._token})) logger.debug("키움 조건검색 LOGIN 발송") except Exception as e: logger.debug("키움 조건검색 LOGIN 발송 실패: %s", e) def _on_error(self, ws, err) -> None: logger.debug("키움 조건검색 WS 오류: %s", err) def _on_close(self, ws, code, msg) -> None: logger.debug("키움 조건검색 WS 종료 (code=%s)", code) def _on_message(self, ws, message) -> None: try: data = json.loads(message) except Exception: return trnm = data.get("trnm") # 키움 WS keep-alive: 받은 PING 을 그대로 돌려보냄 if trnm == "PING": try: ws.send(message) except Exception: pass return if trnm == "LOGIN": if str(data.get("return_code")) in ("0", "0.0"): logger.debug("키움 조건검색 LOGIN OK → CNSRLST") try: ws.send(json.dumps({"trnm": "CNSRLST"})) except Exception: pass else: logger.warning("키움 조건검색 LOGIN 실패: %s", data.get("return_msg")) return if trnm == "CNSRLST": self._handle_condition_list(ws, data.get("data") or []) return if trnm == "CNSRREQ": self._handle_cnsrreq(data) return if trnm == "REAL": self._handle_real(data.get("data") or []) return if trnm == "CNSRCLR": logger.debug("키움 조건검색 CNSRCLR 응답: rc=%s", data.get("return_code")) return # ------------------------------------------------------------------ # 조건식 목록 → seq 해결 → 실시간 등록 # ------------------------------------------------------------------ @staticmethod def _seq_name(item): """CNSRLST data 항목: [seq, name] 배열 또는 {seq,name} dict 모두 허용.""" if isinstance(item, (list, tuple)): seq = str(item[0]) if len(item) > 0 else "" name = str(item[1]) if len(item) > 1 else "" return seq.strip(), name.strip() if isinstance(item, dict): return str(item.get("seq") or "").strip(), str(item.get("name") or "").strip() return "", "" def _handle_condition_list(self, ws, rows: List) -> None: """CNSRLST 응답으로 name→seq 해결 후, 전략별 seq 확정 + CNSRREQ 발송.""" name_to_seq: Dict[str, str] = {} for it in rows: seq, name = self._seq_name(it) if name: name_to_seq[name.strip().lower()] = seq if rows: logger.info( "키움 저장 조건식 %d개: %s", len(rows), ", ".join( f"{self._seq_name(x)[0]}:{self._seq_name(x)[1] or '?'}" for x in rows ), ) # 전략별 seq 확정 (seq 우선, 없으면 name 으로 해결) sid_by_seq: Dict[str, List[str]] = {} for cfg in self._configs: sid = cfg["strategy_id"] seq = (cfg.get("seq") or "").strip() if not seq and cfg.get("name"): seq = name_to_seq.get(cfg["name"].strip().lower(), "") if not seq: logger.warning( "⚠️ 키움 조건식 seq 해결 실패 (strategy=%s name=%s) → 이 전략 폴백", sid, cfg.get("name"), ) continue cfg["seq"] = seq # 해결 결과 반영 sid_by_seq.setdefault(seq, []) if sid not in sid_by_seq[seq]: sid_by_seq[seq].append(sid) logger.info( "🔗 키움 조건식 매핑: strategy=%s seq=%s name=%s", sid, seq, cfg.get("name") or "?", ) with self._kw_lock: self._sid_by_seq = sid_by_seq self._active_seqs = list(sid_by_seq.keys()) # 실시간(search_type=1) 등록 — seq 별 순차 발송 (레이트리밋·응답 누락 방지) self._send_cnsrreq_all(ws) def _send_one_cnsrreq(self, ws, seq: str) -> bool: """단일 seq CNSRREQ 발송.""" payload = { "trnm": "CNSRREQ", "seq": seq, "search_type": "1", "stex_tp": "K", } try: if self._shared_mode and self._shared_ws: return bool(self._shared_ws.send_json(payload)) ws.send(json.dumps(payload)) return True except Exception as e: logger.warning("키움 CNSRREQ 발송 실패 (seq=%s): %s", seq, e) return False def _send_cnsrreq_all(self, ws) -> None: """active seq 목록을 간격 두고 순차 CNSRREQ — 미응답 seq 는 타이머로 재발송.""" with self._kw_lock: seqs = list(self._active_seqs) try: seqs.sort(key=lambda x: int(x)) except Exception: pass gap_lo = float(get_env_float("KIWOOM_CNSRREQ_GAP_MIN_SEC", 0.8)) gap_hi = float(get_env_float("KIWOOM_CNSRREQ_GAP_MAX_SEC", 1.5)) if gap_hi < gap_lo: gap_lo, gap_hi = gap_hi, gap_lo with self._kw_lock: self._cnsrreq_pending = set(seqs) self._cnsrreq_confirmed.clear() sent = 0 for i, seq in enumerate(seqs): if i > 0: time.sleep(random.uniform(gap_lo, gap_hi)) if self._send_one_cnsrreq(ws, seq): sent += 1 logger.debug("키움 CNSRREQ 발송 (seq=%s)", seq) self._start_ok = sent > 0 self._ready.set() self._schedule_cnsrreq_retry(ws, attempt=1) def _schedule_cnsrreq_retry(self, ws, *, attempt: int) -> None: """CNSRREQ 응답이 안 온 seq 만 간격 두고 재발송.""" max_retries = get_env_int("KIWOOM_CNSRREQ_MAX_RETRIES", 3) retry_delay = float(get_env_float("KIWOOM_CNSRREQ_RETRY_SEC", 5.0)) gap_lo = float(get_env_float("KIWOOM_CNSRREQ_GAP_MIN_SEC", 0.8)) gap_hi = float(get_env_float("KIWOOM_CNSRREQ_GAP_MAX_SEC", 1.5)) if gap_hi < gap_lo: gap_lo, gap_hi = gap_hi, gap_lo if self._cnsrreq_retry_timer: try: self._cnsrreq_retry_timer.cancel() except Exception: pass self._cnsrreq_retry_timer = None def _retry() -> None: if not self._running: return with self._kw_lock: missing = sorted( self._cnsrreq_pending - self._cnsrreq_confirmed, key=lambda x: int(x) if str(x).isdigit() else 0, ) if not missing: return if attempt > max_retries: logger.warning( "⚠️ 키움 CNSRREQ 미응답 seq=%s — 최대 재시도 초과", missing, ) return logger.warning( "⚠️ 키움 CNSRREQ 미응답 seq=%s → %.0fs 후 재발송 (%d/%d)", missing, retry_delay, attempt, max_retries, ) for i, seq in enumerate(missing): if i > 0: time.sleep(random.uniform(gap_lo, gap_hi)) self._send_one_cnsrreq(ws, seq) self._schedule_cnsrreq_retry(ws, attempt=attempt + 1) self._cnsrreq_retry_timer = threading.Timer(retry_delay, _retry) self._cnsrreq_retry_timer.daemon = True self._cnsrreq_retry_timer.start() def _handle_cnsrreq(self, data: Dict) -> None: """CNSRREQ 초기 응답: 현재 매칭 종목 리스트로 seq universe 초기화.""" rc = str(data.get("return_code")) seq = str(data.get("seq") or "").strip() if rc not in ("0", "0.0"): logger.warning( "키움 CNSRREQ 실패 (seq=%s) rc=%s msg=%s", seq, rc, data.get("return_msg"), ) return # seq 가 응답에 없을 수 있음 → active_seqs 가 1개면 그걸로 간주 if not seq: with self._kw_lock: if len(self._active_seqs) == 1: seq = self._active_seqs[0] if not seq: return codes: "Dict[str, str]" = {} for it in (data.get("data") or []): code = self._extract_code(it) if code: codes[code] = code # 초기 응답엔 종목명 없음 → code 로 대체 with self._kw_lock: self._seq_codes[seq] = codes self._cnsrreq_confirmed.add(seq) self._cnsrreq_pending.discard(seq) self._publish_seq(seq) logger.info( "✅ 키움 실시간 등록 (seq=%s) 초기 매칭 %d종목", seq, len(codes) ) def _handle_real(self, rows: List) -> None: """REAL push: 843=I(편입)/D(이탈), 9001=종목코드. seq 별 universe 갱신.""" touched: Set[str] = set() with self._kw_lock: for it in rows: vals = it.get("values") if isinstance(it, dict) else None if not isinstance(vals, dict): continue code = _normalize_code(vals.get("9001")) if not code: continue ins_del = str(vals.get("843") or "").strip().upper() # seq 는 item 또는 상위에서 옴. item 에 없으면 유일 seq 로 폴백. seq = str(it.get("item") or it.get("seq") or "").strip() if isinstance(it, dict) else "" if not seq: if len(self._active_seqs) == 1: seq = self._active_seqs[0] else: # seq 불명 + 다중 조건 → 어느 유니버스인지 특정 불가, skip continue bucket = self._seq_codes.setdefault(seq, {}) if ins_del == "D": bucket.pop(code, None) else: # "I" 또는 기타 → 편입으로 처리 (삽입순 유지) if code not in bucket: bucket[code] = code touched.add(seq) for seq in touched: self._publish_seq(seq) @staticmethod def _extract_code(item) -> str: """조건검색 결과 항목에서 종목코드 추출 (9001 우선, jmcode 폴백).""" if isinstance(item, dict): return _normalize_code(item.get("9001") or item.get("jmcode")) if isinstance(item, (list, tuple)) and item: return _normalize_code(item[0]) return "" # ------------------------------------------------------------------ # seq universe → 전략별 반영 (부모 _apply_result 재사용) # ------------------------------------------------------------------ def _publish_seq(self, seq: str) -> None: """seq 의 현재 종목집합을 그 seq 를 쓰는 모든 전략에 반영. 부모 ``_apply_result(strategy_id, rows)`` 를 그대로 호출 → enters/exits 계산·순서·EXIT grace·스냅샷 저장까지 KIS 와 동일하게 처리. """ with self._kw_lock: bucket = dict(self._seq_codes.get(seq, {})) sids = list(self._sid_by_seq.get(seq, [])) # rows: 삽입순(HTS 응답 순서) 유지 → _build_ordered_universe 가 그대로 사용 rows = [{"code": c, "name": n} for c, n in bucket.items()] from ..utils.universe_source import universe_source_active for sid in sids: if not universe_source_active(sid, "kiwoom_condition"): continue try: self._apply_result(sid, rows) except Exception as e: logger.debug("키움 조건검색 _apply_result 예외 (%s): %s", sid, e)