#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ KIS 2번째 앱키(호가 전용 slot=ob) 세션에서 H0STASP0만 41→42 구독 실측. 질문: 호가 전용 키가 세션당 41종까지 되는가? (42번째는 거부인가?) 정책: - REST Approval 재발급 없음 (파일 캐시만). - 봇이 같은 OB 키 WS를 쓰면 테스트 연결이 그 세션을 뺏을 수 있음. 끝나면 close → 봇이 재연결. 메인(tick) 앱키는 사용하지 않음. 사용: cd /home/hoon/kis_bot python3 -u kis_trader/scripts/test_kis_ob_slot_41.py 로그: logs/test_kis_ob_slot_41.log """ from __future__ import annotations import json import logging import sys import threading import time from pathlib import Path from typing import Any, Dict, List, Optional, Tuple HERE = Path(__file__).resolve() ROOT = HERE.parents[2] if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) from kis_approval_manager import KISApprovalManager # noqa: E402 from kis_trader.utils.env import get_env_float, get_env_from_db # noqa: E402 LOG_DIR = ROOT / "logs" LOG_DIR.mkdir(exist_ok=True) LOG_PATH = LOG_DIR / "test_kis_ob_slot_41.log" # 41종 + 42번째 CODES_41 = [ "005930", "000660", "005380", "035420", "051910", "006400", "035720", "068270", "207940", "005490", "028260", "012330", "066570", "003670", "096770", "034730", "015760", "032830", "086790", "009150", "011200", "010130", "018260", "009830", "010950", "024110", "030200", "034220", "047810", "267250", "003550", "017670", "036570", "251270", "352820", "259960", "326030", "377300", "042700", "000270", "105560", # 41 ] CODE_42 = "055550" TR_OB = "H0STASP0" def _setup_logging() -> logging.Logger: lg = logging.getLogger("kis_ob_slot_41") lg.setLevel(logging.DEBUG) lg.propagate = False lg.handlers.clear() fmt = logging.Formatter("[%(asctime)s] %(message)s", datefmt="%H:%M:%S") for h in (logging.StreamHandler(sys.stdout), logging.FileHandler(LOG_PATH, encoding="utf-8")): h.setFormatter(fmt) lg.addHandler(h) return lg def _is_success(rt: str, msg: str) -> bool: rt = (rt or "").strip() msg_u = (msg or "").upper() if rt == "0": return True if "SUBSCRIBE SUCCESS" in msg_u: return True if "ALREADY IN SUBSCRIBE" in msg_u: return True return False class Ob41Probe: def __init__(self, ws_url: str, approval_key: str, logger: logging.Logger, gap_sec: float): self.ws_url = ws_url self.approval_key = approval_key self.logger = logger self.gap_sec = max(0.0, float(gap_sec)) self.acks: List[Dict[str, Any]] = [] self.obs = 0 self.connected_at = 0.0 self.error = "" self.fatal_msg = "" self._ws = None self._send_lock = threading.Lock() self._open_evt = threading.Event() self._pending: List[str] = [] self._pending_lock = threading.Lock() self._alive = True self._ok_codes: List[str] = [] def _build_sub(self, code: str, subscribe: bool = True) -> str: return json.dumps({ "header": { "approval_key": self.approval_key, "custtype": "P", "tr_type": "1" if subscribe else "2", "content-type": "utf-8", }, "body": { "input": { "tr_id": TR_OB, "tr_key": code, } }, }) def send_sub(self, code: str) -> bool: if not self._alive: return False with self._pending_lock: self._pending.append(code) with self._send_lock: if not self._ws or not self._alive: return False try: self._ws.send(self._build_sub(code)) except Exception as e: self._alive = False self.logger.warning("send fail %s: %s", code, e) return False self.logger.info("SENT %s %s", TR_OB, code) return True def _pop_pending(self, code: str) -> str: with self._pending_lock: if code: for i, pc in enumerate(self._pending): if pc == code: return self._pending.pop(i) if self._pending: return self._pending.pop(0) return code def _on_open(self, ws) -> None: self.connected_at = time.time() self.logger.info("on_open OK") self._open_evt.set() def _on_message(self, ws, message: str) -> None: raw = (message or "").strip() if raw == "PINGPONG": try: ws.send("PINGPONG") except Exception: pass return if raw.startswith("{"): try: j = json.loads(raw) except Exception as e: self.logger.warning("JSON parse fail: %s | %s", e, raw[:180]) return hdr = j.get("header") or {} body = j.get("body") or {} out = body.get("output") or {} if not isinstance(out, dict): out = {} tr_key = str(out.get("tr_key") or hdr.get("tr_key") or "").strip() rt = str(body.get("rt_cd") or "").strip() msg1 = str(body.get("msg1") or "").strip() matched = self._pop_pending(tr_key) ok = _is_success(rt, msg1) msg_u = msg1.upper() if "ALREADY IN USE" in msg_u: self.fatal_msg = msg1 self._alive = False row = {"code": matched or tr_key, "ok": ok, "rt_cd": rt, "msg1": msg1} self.acks.append(row) if ok: self._ok_codes.append(row["code"]) self.logger.info("ACK OK code=%s rt=%s %s", row["code"], rt, msg1) else: self.logger.info("ACK REJECT code=%s rt=%s %s", row["code"], rt, msg1) return # pipe 실시간 호가 if "|" in raw and TR_OB in raw: self.obs += 1 def _on_error(self, ws, error) -> None: self.error = str(error) self.logger.warning("on_error: %s", error) def _on_close(self, ws, close_status_code, close_msg) -> None: self._alive = False self.logger.info("on_close code=%s msg=%s", close_status_code, close_msg or "-") def connect(self) -> bool: import websocket 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, ) t = threading.Thread(target=self._ws.run_forever, kwargs={"ping_interval": 20}, daemon=True) t.start() if not self._open_evt.wait(15.0): self.logger.error("on_open timeout") return False return True def close(self) -> None: # 성공 구독 해제 (봇 재연결 전 슬롯 정리) for code in list(self._ok_codes): try: if self._ws and self._alive: with self._send_lock: self._ws.send(self._build_sub(code, subscribe=False)) time.sleep(0.05) except Exception: pass try: if self._ws: self._ws.close() except Exception: pass self._alive = False def _ack_for(acks: List[Dict[str, Any]], code: str) -> Optional[Dict[str, Any]]: for a in reversed(acks): if a.get("code") == code: return a return None def main() -> int: log = _setup_logging() log.info("=" * 64) log.info("KIS OB slot 호가전용 41 한도 실측 → %s", LOG_PATH) log.info("=" * 64) if len(CODES_41) != 41: log.error("CODES_41 길이 %d ≠ 41", len(CODES_41)) return 1 app_key = (get_env_from_db("KIS_APP_KEY_OB_REAL", "") or "").strip() app_secret = (get_env_from_db("KIS_APP_SECRET_OB_REAL", "") or "").strip() if not app_key or not app_secret: log.error("KIS_APP_KEY_OB_REAL / SECRET 없음") return 1 ws_url = (get_env_from_db("KIS_WS_URL_REAL", "") or "").strip() or "ws://ops.koreainvestment.com:21000" gap = float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MIN_SEC", 0.12)) if gap < 0.05: gap = 0.12 mgr = KISApprovalManager.instance(False, slot="ob") key = mgr.reload_from_file() if not key: log.error("approval 파일 캐시 없음 (%s). REST 재발급 금지 → 중단.", mgr._cache_path.name) return 1 log.info("키=KIS_APP_KEY_OB_REAL(앞8자 %s…) slot=ob", app_key[:8]) log.info("WS=%s", ws_url) log.info("approval 캐시 재사용 앞8자 %s… age=%.0f분 REST없음", key[:8], mgr.age_sec() / 60.0) log.info("구독간격=%.2fs | H0STASP0 x 41 + 42=%s", gap, CODE_42) log.info("주의: 봇 kis_ws_ob 가 있으면 세션 뺏김 → 테스트 후 봇 재연결") probe = Ob41Probe(ws_url, key, log, gap) if not probe.connect(): log.error("WS 연결 실패") probe.close() return 1 ok_n = 0 fail_n = 0 first_fail: Optional[Dict[str, Any]] = None log.info("--- Phase A: H0STASP0 x 41 ---") for i, code in enumerate(CODES_41): if not probe._alive: log.warning("세션 종료 — 중단 (%d/41)", i) break if i > 0 and gap > 0: time.sleep(gap) before = len(probe.acks) if not probe.send_sub(code): break # 짧은 ACK 대기 (다음 SENT 전) t0 = time.time() while time.time() - t0 < 2.0 and len(probe.acks) <= before: if not probe._alive: break time.sleep(0.05) a = _ack_for(probe.acks, code) if a and a.get("ok"): ok_n += 1 elif a: fail_n += 1 if first_fail is None: first_fail = a # MAX OVER 면 더 보내지 않음 if "MAX SUBSCRIBE" in str(a.get("msg1") or "").upper(): log.warning("MAX SUBSCRIBE OVER at #%d %s — Phase A 중단", i + 1, code) break if "ALREADY IN USE" in str(a.get("msg1") or "").upper(): break a42: Optional[Dict[str, Any]] = None if probe._alive and ok_n >= 1: log.info("--- Phase B: 42번째 H0STASP0 %s ---", CODE_42) time.sleep(gap) before = len(probe.acks) probe.send_sub(CODE_42) t0 = time.time() while time.time() - t0 < 3.0 and len(probe.acks) <= before: if not probe._alive: break time.sleep(0.05) a42 = _ack_for(probe.acks, CODE_42) time.sleep(1.0) else: log.warning("Phase B 스킵") log.info("연결유지 %.1fs | obs=%d acks=%d err=%s", (time.time() - probe.connected_at) if probe.connected_at else 0.0, probe.obs, len(probe.acks), probe.error or "-") probe.close() def _fmt(a: Optional[Dict[str, Any]]) -> str: if a is None: return "ACK없음" return f"{'OK' if a.get('ok') else 'REJECT'} rt={a.get('rt_cd')} {a.get('msg1')}" log.info("=" * 64) log.info("요약") log.info(" 호가 1~41 OK=%d FAIL=%d (보낸 뒤 첫 MAX/거절=%s)", ok_n, fail_n, _fmt(first_fail) if first_fail else "-") log.info(" 호가 42(%s): %s", CODE_42, _fmt(a42)) log.info(" 실시간 호가 수신 obs=%d", probe.obs) already = "ALREADY IN USE" in (probe.fatal_msg or "").upper() or any( "ALREADY IN USE" in str(a.get("msg1") or "").upper() for a in probe.acks ) if already and ok_n == 0: v = "한도 미측정 — OB 앱키가 봇 WS에 점유(ALREADY IN USE). 앱키당 세션 1개." elif ok_n >= 41 and a42 is not None and not a42.get("ok") and "MAX" in str(a42.get("msg1") or "").upper(): v = "YES — OB 키 호가 41 OK, 42번째는 MAX SUBSCRIBE OVER (세션 한도 41 확인)" elif ok_n >= 41 and a42 is not None and a42.get("ok"): v = "41+42 모두 OK — 서버 한도가 41보다 크거나 ACK 오분류. 로그 재확인" elif ok_n >= 41 and a42 is None: v = "41 OK · 42 ACK없음 — 한도는 최소 41까지는 됨" elif 1 <= ok_n < 41: v = ( f"NO(현재) — OB 키에서 호가 OK={ok_n}/41 만 되고 MAX OVER. " "세션이 비어 있지 않거나(잔여/점유), 앱키·계좌 합산 한도 가능성. " "메인 tick 점유와 합쳐 보면 공유 한도 의심." ) else: v = f"애매 — ok_n={ok_n} fail_n={fail_n}" log.info("판정: %s", v) log.info("로그: %s", LOG_PATH) return 0 if __name__ == "__main__": raise SystemExit(main())