From d737deb47c3b746104190bc793c328c8678c22ce Mon Sep 17 00:00:00 2001 From: Your Name Date: Fri, 28 Aug 2026 16:46:37 +0900 Subject: [PATCH] =?UTF-8?q?fix(ws):=20KIS/=ED=82=A4=EC=9B=80=20WS=C2=B7?= =?UTF-8?q?=EC=A1=B0=EA=B1=B4=EA=B2=80=EC=83=89=20=EC=95=88=EC=A0=95?= =?UTF-8?q?=ED=99=94=20=EB=B0=8F=20=EB=A7=A4=EB=8F=84=EC=8B=9C=EC=84=B8=20?= =?UTF-8?q?=ED=8F=B4=EB=B0=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WS 매니저 spill·조건검색 CNSR 회복·kiwoom_ws_diag 진단 경로를 보강한다. live_sell_price stale RAM/REST 폴백 정합을 유지한다. Co-authored-by: Cursor --- kis_trader/engine/live_sell_price.py | 9 +- .../network/kiwoom_condition_manager.py | 29 +++++-- kis_trader/network/ws_manager.py | 77 +++++++++++++++-- kis_trader/ws/kis_ws.py | 39 ++++++++- kis_trader/ws/kiwoom_ws.py | 44 +++++++++- kis_trader/ws/kiwoom_ws_diag.py | 85 +++++++++++++++++++ 6 files changed, 263 insertions(+), 20 deletions(-) create mode 100644 kis_trader/ws/kiwoom_ws_diag.py diff --git a/kis_trader/engine/live_sell_price.py b/kis_trader/engine/live_sell_price.py index bb9a8ce..26a522a 100644 --- a/kis_trader/engine/live_sell_price.py +++ b/kis_trader/engine/live_sell_price.py @@ -12,7 +12,7 @@ import threading import time from typing import Any, Callable, Optional, Tuple -from kis_trader.utils.env import get_env_float +from kis_trader.utils.env import get_env_bool, get_env_float _stale_rest_lock = threading.Lock() @@ -196,6 +196,7 @@ def resolve_live_sell_price( is_eod: bool = False, fallback_price: float = 0.0, logger: Any = None, + allow_kiwoom_rest: Optional[bool] = None, ) -> Tuple[float, str]: """ 실매 매도용 현재가. @@ -210,6 +211,8 @@ def resolve_live_sell_price( cooldown = float(get_env_float("SELL_WS_STALE_REST_COOLDOWN_SEC", 30.0) or 0.0) if cooldown < 1.0: cooldown = 1.0 + if allow_kiwoom_rest is None: + allow_kiwoom_rest = bool(get_env_bool("SELL_SCAN_ALLOW_KIWOOM_REST", False)) _ = inquire_price # 한투 경로 사용 금지 (시그니처 유지) if is_eod: @@ -217,7 +220,7 @@ def resolve_live_sell_price( if ws_px > 0: _remember_last_good(code, ws_px, "WS") return ws_px, "WS" - if stale_sec > 0 and _stale_rest_allowed(code, cooldown): + if allow_kiwoom_rest and stale_sec > 0 and _stale_rest_allowed(code, cooldown): rest_px = _kiwoom_rest_once(ws, code) if rest_px > 0: return rest_px, "kiwoom_rest" @@ -276,7 +279,7 @@ def resolve_live_sell_price( ) return cached_px, cached_src + "_cached" - if stale_sec > 0 and _stale_rest_allowed(code, cooldown): + if allow_kiwoom_rest and stale_sec > 0 and _stale_rest_allowed(code, cooldown): rest_px = _kiwoom_rest_once(ws, code) if rest_px > 0: _remember_last_good(code, rest_px, "kiwoom_rest") diff --git a/kis_trader/network/kiwoom_condition_manager.py b/kis_trader/network/kiwoom_condition_manager.py index 1e0c22d..fad2b96 100644 --- a/kis_trader/network/kiwoom_condition_manager.py +++ b/kis_trader/network/kiwoom_condition_manager.py @@ -49,6 +49,7 @@ from typing import Any, Dict, List, Optional, Set from .condition_manager import ConditionSearchManager from ..utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int from ..utils.logger import get_logger +from ..ws.kiwoom_ws_diag import log_kiwoom_api_msg, log_kiwoom_ws from ..ws.kis_ws import _get_kiwoom_token_cached logger = get_logger("kis_trader.kwcond") @@ -350,19 +351,26 @@ class KiwoomConditionSearchManager(ConditionSearchManager): ev.set() def _on_shared_trnm(self, ws, msg: dict) -> None: - trnm = msg.get("trnm") + trnm = str(msg.get("trnm") or "") + log_kiwoom_api_msg(logger, trnm or "?", msg, prefix="조건검색") if trnm == "CNSRLST": self._handle_condition_list(ws, msg.get("data") or []) elif trnm == "CNSRREQ": self._handle_cnsrreq(msg) elif trnm == "CNSRCLR": seq_val = _extract_cnsr_seq(msg) - logger.debug( - "키움 조건검색 CNSRCLR 응답: seq=%s rc=%s msg=%s", - seq_val, - msg.get("return_code"), - msg.get("return_msg"), - ) + rc = msg.get("return_code") + rm = msg.get("return_msg") + if str(rc or "").strip() not in ("0", "0.0", ""): + logger.warning( + "키움 CNSRCLR 실패 seq=%s rc=%s msg=%s", + seq_val, rc, rm, + ) + else: + logger.info( + "키움 CNSRCLR OK seq=%s rc=%s msg=%s", + seq_val, rc, rm, + ) self._trigger_cnsrclr_ack(seq_val) def _on_shared_real(self, ws, msg: dict) -> None: @@ -592,6 +600,13 @@ class KiwoomConditionSearchManager(ConditionSearchManager): ws.send(json.dumps(payload)) return True except Exception as e: + log_kiwoom_ws( + logger, "cnsrreq_send_fail", + level="warning", + trnm="CNSRREQ", + extra=f"seq={seq} err={e}", + force=True, + ) logger.warning("키움 CNSRREQ 발송 실패 (seq=%s): %s", seq, e) return False diff --git a/kis_trader/network/ws_manager.py b/kis_trader/network/ws_manager.py index 66e7c80..cbc56a9 100644 --- a/kis_trader/network/ws_manager.py +++ b/kis_trader/network/ws_manager.py @@ -134,6 +134,11 @@ class WSManager: # grace 1회 소진 후 재연장 방지 (재진입 시 discard) self._grace_exhausted: Set[str] = set() self._lock = threading.Lock() + # split reconcile 본문 — WS-ReconcileWorker 단독 (전략 루프는 lock 안 잡음) + self._reconcile_lock = threading.Lock() + self._reconcile_event = threading.Event() + self._reconcile_worker_thread: Optional[threading.Thread] = None + self._reconcile_ls_feed_pending = False # 갭보정 WS 재접속 시: split 모드면 KIS∪키움 관심 종목 전체 self._gap_refill_codes: Set[str] = set() # BaseStrategy 틱매도 등 — 시세 캐시 갱신 리스너 @@ -335,6 +340,7 @@ class WSManager: self._start_share_meta_worker() self._start_gap_worker() self._start_ls_gap_worker() + self._start_reconcile_worker() # 영구 구독 (KOSPI/KOSDAQ ETF 등) — KR 은 LS WS. KIS 슬롯·ws_candles 갭 제외 self._load_permanent_codes() @@ -463,8 +469,8 @@ class WSManager: with self._lock: self._owner_candidates[owner] = cand self._owner_holdings[owner] = hold - self._reconcile_split_subscriptions() - self._reconcile_ls_feed_subscriptions() + # REG/호가/갭 enqueue 는 전략 루프 밖 단일 워커 (동시 4전략 lock 대기 금지) + self._schedule_reconcile_split(include_ls_feed=True) def _ls_feed_codes_locked(self) -> Set[str]: """호출자 _lock 보유 가정.""" @@ -1287,13 +1293,69 @@ class WSManager: self._ob_home.pop(code, None) self._ob_home_spill.discard(code) - def _reconcile_split_subscriptions(self) -> None: - """MINIMAL 재동기화: 키움=후보∪보유, 한투=보유만, 영구KR=LS (한투/키움 슬롯 제외). + def _start_reconcile_worker(self) -> None: + """split 구독 reconcile 전담 — 갭 워커와 동일 패턴.""" + if self._reconcile_worker_thread and self._reconcile_worker_thread.is_alive(): + return + t = threading.Thread( + target=self._reconcile_worker_loop, + name="WS-ReconcileWorker", + daemon=True, + ) + t.start() + self._reconcile_worker_thread = t + logger.info("✅ WS reconcile 워커 시작 (전략 루프 논블로킹)") - LIVE_TICK 을 읽지 않음. 메인 벤더와 구독 명단은 별개. - """ + def _schedule_reconcile_split(self, *, include_ls_feed: bool = False) -> None: + """owner 후보/보유 RAM 갱신 후 reconcile 요청 — debounce 로 4전략 합침.""" if not (self._split_feed_active and self.ws_cache and self._kiwoom_ws): return + if include_ls_feed: + self._reconcile_ls_feed_pending = True + self._reconcile_event.set() + self._start_reconcile_worker() + + def _reconcile_worker_loop(self) -> None: + """debounce 후 1회 reconcile — 전략 쓰레드는 여기서 블록되지 않음.""" + while True: + self._reconcile_event.wait() + debounce = float(get_env_float("WS_RECONCILE_DEBOUNCE_SEC", 0.25) or 0.0) + if debounce > 0: + time.sleep(debounce) + while self._reconcile_event.is_set(): + self._reconcile_event.clear() + time.sleep(debounce) + else: + self._reconcile_event.clear() + + if not (self._split_feed_active and self.ws_cache and self._kiwoom_ws): + self._reconcile_ls_feed_pending = False + continue + + t0 = time.perf_counter() + with self._reconcile_lock: + try: + self._reconcile_split_subscriptions_locked() + except Exception as ex: + logger.warning("⚠️ [reconcile worker] split 예외: %s", ex) + split_ms = (time.perf_counter() - t0) * 1000.0 + + run_ls = self._reconcile_ls_feed_pending + self._reconcile_ls_feed_pending = False + if run_ls: + try: + self._reconcile_ls_feed_subscriptions() + except Exception as ex: + logger.debug("reconcile ls_feed 예외: %s", ex) + + if split_ms >= 3000.0: + logger.warning( + "⚠️ [reconcile worker] split %.0fms (REG/호가/갭 enqueue)", + split_ms, + ) + + def _reconcile_split_subscriptions_locked(self) -> None: + """_reconcile_lock 보유 가정 — 본문.""" # permanent_subscriptions 테이블 갱신 반영 (보유 해제 후에도 영구구독 틱 유지) now = time.time() if now - self._permanent_reload_ts >= 300.0: @@ -1795,9 +1857,10 @@ class WSManager: for code in added_hold + added_rest: self.subscribe(code, owner) self._sync_tick_record_codes() + holdings = self._all_holdings() with self._lock: perm = set(self._permanent_codes) - kw_want = set(self._code_refs.keys()) | self._all_holdings() + kw_want = set(self._code_refs.keys()) | holdings self._finalize_ls_subscriptions(kw_want, perm) def get_recent_ticks(self, code: str, limit: int = 100) -> list: diff --git a/kis_trader/ws/kis_ws.py b/kis_trader/ws/kis_ws.py index 1247937..67b9833 100644 --- a/kis_trader/ws/kis_ws.py +++ b/kis_trader/ws/kis_ws.py @@ -2614,6 +2614,39 @@ def _ka10080_semaphore() -> _threading.Semaphore: return _ka10080_sem +def _ka10080_sem_timeout_sec(*, purpose: str = "") -> float: + """ka10080/ka10007 공유 sem 대기 상한(초). 0=무제한(레거시).""" + key = "KIWOOM_KA10007_SEM_TIMEOUT_SEC" if purpose == "ka10007" else "KIWOOM_KA10080_SEM_TIMEOUT_SEC" + default = 2.0 if purpose == "ka10007" else 12.0 + try: + return max(0.0, float(get_env_float(key, default) or 0.0)) + except Exception: + return default + + +def _ka10080_acquire( + sem: _threading.Semaphore, + *, + purpose: str, + code: str = "", +) -> bool: + """갭보정(ka10080)·매도 REST(ka10007) sem — 무한 대기 금지.""" + wait = _ka10080_sem_timeout_sec(purpose=purpose) + if wait <= 0: + sem.acquire() + return True + if sem.acquire(timeout=wait): + return True + logger.warning( + "⚠️ 키움 REST sem 타임아웃 purpose=%s code=%s wait=%.0fs " + "(갭 ka10080 점유 중 — 매도·매수 루프 기아 방지)", + purpose, + code or "-", + wait, + ) + return False + + def _ka10080_is_rate_limit(http_st: int, rc: Any, msg: str) -> bool: if int(http_st or 0) == 429: return True @@ -2857,7 +2890,8 @@ def get_kiwoom_candles_df( data = None resp = None for attempt in range(rate_retries): - sem.acquire() + if not _ka10080_acquire(sem, purpose="ka10080", code=code): + return pd.DataFrame() try: resp = requests.post(base_url, json=body, headers=headers, timeout=15) http_st = int(resp.status_code) @@ -2978,7 +3012,8 @@ def fetch_kiwoom_cur_prc_ka10007( sem = _ka10080_semaphore() data: Dict[str, Any] = {} http_st = 0 - sem.acquire() + if not _ka10080_acquire(sem, purpose="ka10007", code=code): + return 0.0 try: resp = requests.post(url, json=body, headers=headers, timeout=8) http_st = int(resp.status_code) diff --git a/kis_trader/ws/kiwoom_ws.py b/kis_trader/ws/kiwoom_ws.py index d84d4a9..15060a6 100644 --- a/kis_trader/ws/kiwoom_ws.py +++ b/kis_trader/ws/kiwoom_ws.py @@ -93,6 +93,8 @@ 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 + # ────────────────────────────────────────────────────────────────────── # 메인 클래스 @@ -364,11 +366,25 @@ class KiwoomWebSocketPriceCache: 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 @@ -667,6 +683,9 @@ class KiwoomWebSocketPriceCache: 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: @@ -728,7 +747,10 @@ class KiwoomWebSocketPriceCache: if trnm == "REG": rc = msg.get("return_code") if rc != 0: - logger.warning("⚠️ 키움 WS REG 실패: %s", msg.get("return_msg", "")) + logger.warning( + "⚠️ 키움 WS REG 실패 rc=%s msg=%s", + rc, msg.get("return_msg", ""), + ) return if trnm == "REMOVE": @@ -955,11 +977,25 @@ class KiwoomWebSocketPriceCache: 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 "") @@ -1079,6 +1115,12 @@ class KiwoomWebSocketPriceCache: 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 diff --git a/kis_trader/ws/kiwoom_ws_diag.py b/kis_trader/ws/kiwoom_ws_diag.py new file mode 100644 index 0000000..cde6bfe --- /dev/null +++ b/kis_trader/ws/kiwoom_ws_diag.py @@ -0,0 +1,85 @@ +"""키움 WS 진단 로그 — return_code/return_msg·스레드·trnm 추적.""" +from __future__ import annotations + +import threading +from typing import Any, Mapping, Optional + +from kis_trader.utils.env import get_env_bool + + +def kiwoom_ws_diag_enabled() -> bool: + try: + return bool(get_env_bool("KIWOOM_WS_DIAG_ENABLED", True)) + except Exception: + return True + + +def _thread_tag() -> str: + t = threading.current_thread() + return f"{t.name}#{getattr(t, 'ident', 0) or 0}" + + +def format_kiwoom_return(msg: Mapping[str, Any]) -> str: + rc = msg.get("return_code") + rm = msg.get("return_msg") + if rc is None and rm in (None, ""): + return "" + return f"rc={rc} msg={rm}" + + +def log_kiwoom_ws( + logger: Any, + event: str, + *, + level: str = "info", + rc: Any = None, + msg: Any = None, + trnm: Optional[str] = None, + extra: Optional[str] = None, + force: bool = False, +) -> None: + """KIWOOM_WS_DIAG_ENABLED 또는 force=True 일 때 INFO/WARNING 로그.""" + if not force and not kiwoom_ws_diag_enabled(): + return + parts = [f"🔎 [키움WS/{event}]", f"thr={_thread_tag()}"] + if trnm: + parts.append(f"trnm={trnm}") + if rc is not None or (msg not in (None, "")): + parts.append(f"rc={rc} msg={msg}") + if extra: + parts.append(str(extra)) + line = " ".join(parts) + try: + getattr(logger, level, logger.info)(line) + except Exception: + pass + + +def log_kiwoom_api_msg( + logger: Any, + trnm: str, + data: Mapping[str, Any], + *, + level: str = "warning", + prefix: str = "", +) -> None: + """수신 JSON에 return_code 가 있고 0이 아니면 항상 로그 (진단 OFF여도).""" + rc_raw = data.get("return_code") + if rc_raw is None: + return + rc = str(rc_raw).strip() + if rc in ("0", "0.0", ""): + return + rm = str(data.get("return_msg") or "").strip() + seq = data.get("seq") or data.get("data") + extra = f"seq={seq}" if seq not in (None, "", []) else "" + head = f"{prefix} " if prefix else "" + line = f"{head}키움 {trnm} rc={rc} msg={rm}" + if extra: + line = f"{line} {extra}" + if kiwoom_ws_diag_enabled(): + line = f"🔎 {line} thr={_thread_tag()}" + try: + getattr(logger, level, logger.warning)(line) + except Exception: + pass