- Optuna web jobs/TPE/apply snapshot·틱로더 정합, jobs limit·감사로그 - 백테 UI 호가모드·후보 적용 흐름, feed_collect_stats API/탭 - 가설검증·교차검증 룰, 4전략 스모크·OB slot41 진단 스크립트 Co-authored-by: Cursor <cursoragent@cursor.com>
374 lines
13 KiB
Python
374 lines
13 KiB
Python
#!/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())
|