feat(옵투나·웹): 후처리 재탐색·ob_modes·적용감사·수집통계

- Optuna web jobs/TPE/apply snapshot·틱로더 정합, jobs limit·감사로그
- 백테 UI 호가모드·후보 적용 흐름, feed_collect_stats API/탭
- 가설검증·교차검증 룰, 4전략 스모크·OB slot41 진단 스크립트

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Your Name
2026-08-27 15:23:44 +09:00
parent 5e44b86f8b
commit 8fbba264ba
30 changed files with 2973 additions and 186 deletions

View File

@@ -0,0 +1,373 @@
#!/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())