Changes: - Introduced new files for strategy definitions and study names. - Enhanced `backtest_web.py` with functions to handle integer display prices and trade data formatting. - Updated backtesting logic to incorporate end-of-day (EOD) parameters for breakout and momentum strategies. - Added EOD configuration options in the database and parameter search files. Impact: - These changes improve the modularity and usability of the backtesting framework, allowing for better integration of EOD strategies and clearer trade data presentation.
841 lines
34 KiB
Python
841 lines
34 KiB
Python
"""
|
|
kis_trader/ws/kiwoom_ws.py — 키움 WebSocket 실시간 시세 캐시 (시세 마이그레이션 검증용)
|
|
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
|
|
|
|
목적
|
|
----
|
|
KIS WS(41 한도)의 시세를 키움 WS(100 한도)로 옮기기 전에, 둘을 동시에 돌려
|
|
가격 일치성을 검증하기 위한 키움 WebSocket 클라이언트.
|
|
|
|
설계 원칙 (KIS WS 와 동일)
|
|
--------------------------
|
|
- ``Dict[str, Dict]`` 메모리 캐시 (락만 잠그고 마이크로초 read/write)
|
|
- ``get_price(code)`` 인터페이스를 KIS WS 와 100% 동일 포맷으로 제공
|
|
→ 봇 코드 재사용성 100%
|
|
- ``KiwoomTokenManager`` 싱글톤 (kis_ws.py 안에 있음) 재사용
|
|
- 재연결 백오프, approval 갱신, 종목 등록/해지
|
|
|
|
키움 WebSocket 스펙
|
|
-------------------
|
|
URL 실전: wss://api.kiwoom.com:10000/api/dostk/websocket
|
|
URL 모의: wss://mockapi.kiwoom.com:10000/api/dostk/websocket
|
|
|
|
[프로토콜]
|
|
1. 연결 후 LOGIN: {"trnm": "LOGIN", "token": "<access_token>"}
|
|
2. 등록: {"trnm": "REG", "grp_no": "1", "refresh": "1",
|
|
"data": [{"item": ["005930"], "type": ["0B"]}]}
|
|
3. 해지: {"trnm": "REMOVE","grp_no": "1",
|
|
"data": [{"item": ["005930"], "type": ["0B"]}]}
|
|
4. PING: {"trnm": "PING"} → 30초마다 서버가 보냄, 그대로 echo
|
|
|
|
[주식체결 0B 메시지 FID]
|
|
10 = 현재가(체결가; 부호 → 음수=하락) 11 = 전일대비
|
|
12 = 등락률 13 = 누적거래량
|
|
14 = 누적거래대금 15 = 거래량(체결량)
|
|
16 = 시가 17 = 고가 18 = 저가 20 = 체결시간(HHMMSS)
|
|
|
|
(키움 OpenAPI+ FID 정의 — 종목별 동일)
|
|
|
|
설치
|
|
----
|
|
``pip install websocket-client`` (KIS WS 와 공용)
|
|
|
|
토글 (DB env_config)
|
|
--------------------
|
|
``WS_PROVIDER`` (기본 ``kis_only`` — 키움 WS 미기동)
|
|
``KIWOOM_WS_URL_REAL`` / ``KIWOOM_WS_URL_MOCK`` (URL 재정의용)
|
|
``KIWOOM_WS_REG_CHUNK_SIZE`` / ``KIWOOM_WS_REG_GAP_SEC`` / ``KIWOOM_WS_REG_DEBOUNCE_SEC``
|
|
— REG(TRNM) 초당 허용 건수 초과 방지(배치·전송 간격·단건 디바운스).
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import threading
|
|
import time
|
|
from collections import defaultdict
|
|
from typing import Any, Callable, Dict, Iterable, List, Optional, Set
|
|
|
|
logger = logging.getLogger("KiwoomWebSocket")
|
|
|
|
|
|
# ── 환경변수 헬퍼 (KIS WS 와 동일 패턴, fallback 포함) ────────────────────
|
|
try:
|
|
from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int # noqa: F401
|
|
except ImportError:
|
|
try:
|
|
from kis_long_ver1 import get_env_float, get_env_from_db, get_env_int # noqa: F401
|
|
except ImportError:
|
|
def get_env_from_db(key, default=""): # type: ignore[misc]
|
|
return default
|
|
|
|
def get_env_int(key, default): # type: ignore[misc]
|
|
return default
|
|
|
|
def get_env_float(key, default): # type: ignore[misc]
|
|
return float(default)
|
|
|
|
def get_env_bool(key, default=False): # type: ignore[misc]
|
|
return default
|
|
|
|
|
|
try:
|
|
from .orderbook_cache import OrderbookCache, OrderbookSnapshot
|
|
except ImportError:
|
|
OrderbookCache = None # type: ignore[misc, assignment]
|
|
OrderbookSnapshot = None # type: ignore[misc, assignment]
|
|
|
|
try:
|
|
from .program_cache import ProgramCache, ProgramSnapshot
|
|
except ImportError:
|
|
ProgramCache = None # type: ignore[misc, assignment]
|
|
ProgramSnapshot = None # type: ignore[misc, assignment]
|
|
|
|
|
|
# ──────────────────────────────────────────────────────────────────────
|
|
# 메인 클래스
|
|
# ──────────────────────────────────────────────────────────────────────
|
|
class KiwoomWebSocketPriceCache:
|
|
"""키움 실시간 체결가 (0B) WebSocket 수신기.
|
|
|
|
사용법::
|
|
|
|
ws = KiwoomWebSocketPriceCache(app_key, app_secret, is_mock=False)
|
|
ws.start()
|
|
ws.subscribe("005930")
|
|
data = ws.get_price("005930") # KIS 와 동일 dict 포맷
|
|
ws.stop()
|
|
|
|
``get_price()`` 가 None 이면 → 캐시 없음/만료 → 호출자가 다른 소스 사용.
|
|
"""
|
|
|
|
# 키움 0B FID
|
|
FID_PRICE = "10" # 현재가 (체결가, 부호 포함)
|
|
FID_CHANGE = "11" # 전일대비
|
|
FID_CHANGE_PCT = "12" # 등락률
|
|
FID_CUM_VOL = "13" # 누적거래량
|
|
FID_OPEN = "16"
|
|
FID_HIGH = "17"
|
|
FID_LOW = "18"
|
|
FID_TICK_VOL = "15"
|
|
FID_TICK_TIME = "20"
|
|
FID_EXEC_STRENGTH = "567" # 체결강도 (%)
|
|
|
|
# 한도 (키움 권장)
|
|
MAX_SUBSCRIPTIONS_PER_GROUP = 100 # grp_no=1 그룹 1개당
|
|
GROUP_NO = "1"
|
|
SUB_TYPE = "0B" # 주식체결
|
|
SUB_TYPE_ORDERBOOK = "0D" # 주식호가잔량
|
|
SUB_TYPE_PROGRAM = "0w" # 종목프로그램매매
|
|
|
|
# 재연결 정책
|
|
RECONNECT_BASE_DELAY_SEC = 5.0
|
|
RECONNECT_MAX_DELAY_SEC = 300.0
|
|
MAX_RECONNECTS_PER_HOUR = 6
|
|
MAX_RECONNECT_ATTEMPTS = 10
|
|
STABLE_CONN_RESET_SEC = 300.0 # 5분 안정 연결 후 끊기면 카운터 초기화
|
|
|
|
# 토큰 캐시 — KiwoomTokenManager 가 알아서 처리하지만 보수적 만료 버퍼
|
|
TOKEN_REFRESH_BUFFER_SEC = 600
|
|
|
|
def __init__(
|
|
self,
|
|
app_key: str,
|
|
app_secret: str,
|
|
is_mock: bool = False,
|
|
):
|
|
self.app_key = app_key
|
|
self.app_secret = app_secret
|
|
self.is_mock = is_mock
|
|
|
|
# URL (env/DB 로 재정의 가능)
|
|
_default_real = "wss://api.kiwoom.com:10000/api/dostk/websocket"
|
|
_default_mock = "wss://mockapi.kiwoom.com:10000/api/dostk/websocket"
|
|
self._ws_url = (
|
|
get_env_from_db("KIWOOM_WS_URL_MOCK", _default_mock)
|
|
if is_mock
|
|
else get_env_from_db("KIWOOM_WS_URL_REAL", _default_real)
|
|
)
|
|
|
|
# 메모리 캐시 — KIS WS 와 동일 포맷 (data + ts)
|
|
self._cache: Dict[str, Dict] = {}
|
|
self._cache_lock = threading.Lock()
|
|
self._orderbook_cache = OrderbookCache() if OrderbookCache else None
|
|
self._program_cache = ProgramCache() if ProgramCache else None
|
|
|
|
# 구독 종목
|
|
self._subscribed: Set[str] = set()
|
|
self._sub_lock = threading.Lock()
|
|
|
|
# 선택: KIS CandleAggregator 에 틱 전달 (후보 종목만 필터링 가능)
|
|
self._candle_agg: Any = None
|
|
# None = 구독 전 종목 틱을 집계기에 전달, Set = 해당 코드만 전달
|
|
self._candle_agg_codes: Optional[Set[str]] = None
|
|
# TickRecorder (C안 ws_ticks) — 필터는 TickRecorder.set_record_codes()
|
|
self._tick_recorder: Any = None
|
|
self._trigger_snapshot_recorder: Any = None
|
|
|
|
# 연결 상태
|
|
self._ws = None
|
|
self._ws_thread: Optional[threading.Thread] = None
|
|
self._running = False
|
|
self._connected = False
|
|
self._authenticated = False # LOGIN 응답 OK 받기 전엔 REG 못 보냄
|
|
|
|
# 재연결 추적
|
|
self._reconnect_count = 0
|
|
self._reconnect_times: list = []
|
|
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
|
|
self._last_connect_time: float = 0.0
|
|
|
|
# REG 레이트리밋 회피: 단건 subscribe 는 디바운스 후 묶어 전송
|
|
self._reg_batch_codes: Set[str] = set()
|
|
self._reg_timer: Optional[threading.Timer] = None
|
|
self._reg_timer_lock = threading.Lock()
|
|
|
|
# 조건검색 등 외부 모듈 — **동일 WS 세션 공유** (키움은 토큰당 1접속)
|
|
self._ext_handler_lock = threading.Lock()
|
|
self._ext_handlers: Dict[str, List[Callable]] = defaultdict(list)
|
|
self._login_callbacks: List[Callable] = []
|
|
self._login_cb_lock = threading.Lock()
|
|
|
|
# websocket-client lib
|
|
try:
|
|
import websocket as _ws_lib # type: ignore
|
|
self._ws_lib = _ws_lib
|
|
self._available = True
|
|
except ImportError:
|
|
self._ws_lib = None
|
|
self._available = False
|
|
logger.warning("⚠️ websocket-client 미설치 — 키움 WS 사용 불가")
|
|
|
|
# ------------------------------------------------------------------
|
|
# 외부 API
|
|
# ------------------------------------------------------------------
|
|
def start(self) -> bool:
|
|
"""백그라운드 수신 스레드 기동."""
|
|
if not self._available:
|
|
logger.warning("키움 WS 라이브러리 없음 → start 무시")
|
|
return False
|
|
if self._running:
|
|
return True
|
|
if not self.app_key or not self.app_secret:
|
|
logger.warning("⚠️ 키움 키 없음 → 키움 WS 비활성")
|
|
return False
|
|
|
|
self._running = True
|
|
self._ws_thread = threading.Thread(
|
|
target=self._run_loop, daemon=True, name="KiwoomWS",
|
|
)
|
|
self._ws_thread.start()
|
|
logger.info("✅ 키움 WebSocket 수신 스레드 시작 (mock=%s, url=%s)",
|
|
self.is_mock, self._ws_url)
|
|
return True
|
|
|
|
def stop(self) -> None:
|
|
"""수신 스레드 종료 + 소켓 닫기."""
|
|
self._running = False
|
|
with self._reg_timer_lock:
|
|
if self._reg_timer:
|
|
try:
|
|
self._reg_timer.cancel()
|
|
except Exception:
|
|
pass
|
|
self._reg_timer = None
|
|
with self._sub_lock:
|
|
self._reg_batch_codes.clear()
|
|
try:
|
|
if self._ws is not None:
|
|
self._ws.close()
|
|
except Exception:
|
|
pass
|
|
|
|
def _max_subscriptions(self) -> int:
|
|
"""그룹당 최대 구독 수 — DB ``KIWOOM_WS_MAX_SUBSCRIPTIONS`` (기본 100)."""
|
|
v = get_env_int("KIWOOM_WS_MAX_SUBSCRIPTIONS", self.MAX_SUBSCRIPTIONS_PER_GROUP)
|
|
return max(1, min(200, int(v)))
|
|
|
|
def _reg_chunk_size(self) -> int:
|
|
"""REG 한 메시지당 최대 종목 수 — ``KIWOOM_WS_REG_CHUNK_SIZE`` (기본 25)."""
|
|
v = get_env_int("KIWOOM_WS_REG_CHUNK_SIZE", 25)
|
|
return max(1, min(80, int(v)))
|
|
|
|
def _reg_gap_sec(self) -> float:
|
|
"""청크 사이 전송 간격(초) — ``KIWOOM_WS_REG_GAP_SEC`` (기본 0.18)."""
|
|
return max(0.05, min(2.0, get_env_float("KIWOOM_WS_REG_GAP_SEC", 0.18)))
|
|
|
|
def _reg_debounce_sec(self) -> float:
|
|
"""단건 subscribe 묶기 대기(초) — ``KIWOOM_WS_REG_DEBOUNCE_SEC`` (기본 0.12)."""
|
|
return max(0.0, min(1.0, get_env_float("KIWOOM_WS_REG_DEBOUNCE_SEC", 0.12)))
|
|
|
|
def _orderbook_ws_enabled(self) -> bool:
|
|
return get_env_bool("KIWOOM_WS_ORDERBOOK_ENABLED", True)
|
|
|
|
def _program_ws_enabled(self) -> bool:
|
|
return get_env_bool("KIWOOM_WS_PROGRAM_ENABLED", True)
|
|
|
|
def is_available(self) -> bool:
|
|
"""websocket-client 설치 및 키 설정 여부."""
|
|
return bool(self._available and self.app_key and self.app_secret)
|
|
|
|
def is_authenticated(self) -> bool:
|
|
"""LOGIN OK 이후 REG/조건검색 전송 가능."""
|
|
return bool(self._connected and self._authenticated)
|
|
|
|
def register_trnm_handler(self, trnm: str, fn: Callable) -> None:
|
|
"""외부 모듈(조건검색 CNSR* 등)용 trnm 핸들러 — 시세 WS 와 세션 공유."""
|
|
key = str(trnm or "").strip().upper()
|
|
if not key or not callable(fn):
|
|
return
|
|
with self._ext_handler_lock:
|
|
if fn not in self._ext_handlers[key]:
|
|
self._ext_handlers[key].append(fn)
|
|
|
|
def unregister_trnm_handler(self, trnm: str, fn: Callable) -> None:
|
|
key = str(trnm or "").strip().upper()
|
|
with self._ext_handler_lock:
|
|
lst = self._ext_handlers.get(key)
|
|
if lst and fn in lst:
|
|
lst.remove(fn)
|
|
|
|
def add_on_login_callback(self, fn: Callable) -> None:
|
|
"""LOGIN OK 직후(재접속마다) 호출 — 조건검색 CNSRLST 등."""
|
|
if not callable(fn):
|
|
return
|
|
with self._login_cb_lock:
|
|
if fn not in self._login_callbacks:
|
|
self._login_callbacks.append(fn)
|
|
|
|
def remove_on_login_callback(self, fn: Callable) -> None:
|
|
with self._login_cb_lock:
|
|
if fn in self._login_callbacks:
|
|
self._login_callbacks.remove(fn)
|
|
|
|
def send_json(self, msg: dict) -> bool:
|
|
"""인증된 WS 에 JSON 전송 (조건검색 CNSRREQ 등)."""
|
|
if not self.is_authenticated() or not self._ws:
|
|
return False
|
|
try:
|
|
self._ws.send(json.dumps(msg))
|
|
return True
|
|
except Exception as e:
|
|
logger.debug("키움 WS send_json 실패: %s", e)
|
|
return False
|
|
|
|
def _fire_login_callbacks(self, ws) -> None:
|
|
with self._login_cb_lock:
|
|
cbs = list(self._login_callbacks)
|
|
for cb in cbs:
|
|
try:
|
|
cb(ws)
|
|
except Exception as e:
|
|
logger.debug("키움 WS login callback 예외: %s", e)
|
|
|
|
def _dispatch_ext_handlers(self, trnm: str, ws, msg: dict) -> None:
|
|
key = str(trnm or "").strip().upper()
|
|
with self._ext_handler_lock:
|
|
handlers = list(self._ext_handlers.get(key, []))
|
|
for fn in handlers:
|
|
try:
|
|
fn(ws, msg)
|
|
except Exception as e:
|
|
logger.debug("키움 WS ext handler(%s) 예외: %s", key, e)
|
|
|
|
def _reg_types(self) -> List[str]:
|
|
"""REG/REMOVE 실시간 타입 — 0B 체결 + (옵션) 0D 호가 + 0w 프로그램."""
|
|
types = [self.SUB_TYPE]
|
|
if self._orderbook_ws_enabled():
|
|
types.append(self.SUB_TYPE_ORDERBOOK)
|
|
if self._program_ws_enabled():
|
|
types.append(self.SUB_TYPE_PROGRAM)
|
|
return types
|
|
|
|
def subscribe_many(self, codes: Iterable[str]) -> List[str]:
|
|
"""여러 종목 등록. REG 는 청크+간격으로 전송(TRNM 레이트리밋 회피).
|
|
|
|
LOGIN 전이면 집합만 갱신하고, LOGIN OK 시 ``_send_reg_chunked`` 로 일괄 전송.
|
|
"""
|
|
added: List[str] = []
|
|
with self._sub_lock:
|
|
for raw in codes:
|
|
c = (str(raw) or "").strip()
|
|
if not c or c in self._subscribed:
|
|
continue
|
|
if len(self._subscribed) >= self._max_subscriptions():
|
|
logger.warning(
|
|
"⚠️ 키움 WS 구독 한도 초과 (%d/%d) — 이후 종목 스킵",
|
|
len(self._subscribed), self._max_subscriptions(),
|
|
)
|
|
break
|
|
self._subscribed.add(c)
|
|
added.append(c)
|
|
if added and self._connected and self._authenticated:
|
|
self._send_reg_chunked(added)
|
|
return added
|
|
|
|
def subscribe(self, code: str) -> bool:
|
|
"""단일 종목 등록. REG 는 짧게 디바운스 후 묶어 전송."""
|
|
code = (code or "").strip()
|
|
if not code:
|
|
return False
|
|
with self._sub_lock:
|
|
if code in self._subscribed:
|
|
return True
|
|
if len(self._subscribed) >= self._max_subscriptions():
|
|
logger.warning(
|
|
"⚠️ 키움 WS 구독 한도 초과 (%d/%d) — %s 등록 거절",
|
|
len(self._subscribed), self._max_subscriptions(), code,
|
|
)
|
|
return False
|
|
self._subscribed.add(code)
|
|
# LOGIN 전에는 집합만 쌓고, REG 는 LOGIN OK 후 일괄 전송(중복·순서 레이스 방지)
|
|
need_reg = self._connected and self._authenticated
|
|
if need_reg:
|
|
self._reg_batch_codes.add(code)
|
|
if need_reg:
|
|
self._schedule_reg_debounce()
|
|
return True
|
|
|
|
def _schedule_reg_debounce(self) -> None:
|
|
deb = self._reg_debounce_sec()
|
|
if deb <= 0:
|
|
with self._sub_lock:
|
|
pending = sorted(self._reg_batch_codes)
|
|
self._reg_batch_codes.clear()
|
|
if pending:
|
|
self._send_reg_chunked(pending)
|
|
return
|
|
with self._reg_timer_lock:
|
|
if self._reg_timer:
|
|
try:
|
|
self._reg_timer.cancel()
|
|
except Exception:
|
|
pass
|
|
self._reg_timer = threading.Timer(deb, self._flush_reg_debounced)
|
|
self._reg_timer.daemon = True
|
|
self._reg_timer.start()
|
|
|
|
def _flush_reg_debounced(self) -> None:
|
|
with self._reg_timer_lock:
|
|
self._reg_timer = None
|
|
with self._sub_lock:
|
|
pending = sorted(self._reg_batch_codes)
|
|
self._reg_batch_codes.clear()
|
|
if pending and self._connected and self._authenticated:
|
|
self._send_reg_chunked(pending)
|
|
|
|
def unsubscribe(self, code: str) -> bool:
|
|
"""단일 종목 등록 해지."""
|
|
code = (code or "").strip()
|
|
with self._sub_lock:
|
|
if code not in self._subscribed:
|
|
return False
|
|
self._subscribed.discard(code)
|
|
if self._orderbook_cache:
|
|
self._orderbook_cache.remove(code)
|
|
if self._program_cache:
|
|
self._program_cache.remove(code)
|
|
if self._connected and self._authenticated:
|
|
return self._send_remove([code])
|
|
return True
|
|
|
|
def get_orderbook_snapshot(
|
|
self, code: str, max_age_sec: float = 3.0,
|
|
) -> Optional["OrderbookSnapshot"]:
|
|
"""키움 0D RAM 호가 스냅샷."""
|
|
if not self._orderbook_cache:
|
|
return None
|
|
return self._orderbook_cache.get(code, max_age_sec=max_age_sec)
|
|
|
|
def get_orderbook(self, code: str, max_age_sec: float = 3.0) -> Optional[Dict]:
|
|
"""KIS REST 호가 dict 호환 — ``orderbook_sell`` / OrderManager 폴백용."""
|
|
if not self._orderbook_cache:
|
|
return None
|
|
return self._orderbook_cache.get_kis_bid_dict(code, max_age_sec=max_age_sec)
|
|
|
|
def get_program_snapshot(
|
|
self, code: str, max_age_sec: float = 30.0,
|
|
) -> Optional["ProgramSnapshot"]:
|
|
"""키움 0w RAM 프로그램매매 스냅샷."""
|
|
if not self._program_cache:
|
|
return None
|
|
return self._program_cache.get(code, max_age_sec=max_age_sec)
|
|
|
|
def get_price(self, code: str, max_age_sec: float = 5.0) -> Optional[Dict]:
|
|
"""KIS WS ``get_price`` 와 동일 포맷 반환.
|
|
|
|
반환::
|
|
|
|
{
|
|
"stck_prpr": "73900", # 현재가
|
|
"stck_oprc": "73000", # 시가
|
|
"stck_hgpr": "74500", # 고가
|
|
"stck_lwpr": "72800", # 저가
|
|
"prdy_vrss": "200", # 전일 대비
|
|
"prdy_ctrt": "0.27", # 등락률
|
|
"_age_ms": 123, # 캐시 나이 (ms) — 검증용 메타
|
|
}
|
|
|
|
``max_age_sec`` 초 초과면 None.
|
|
"""
|
|
with self._cache_lock:
|
|
entry = self._cache.get(code)
|
|
if not entry:
|
|
return None
|
|
age_sec = time.time() - entry.get("ts", 0)
|
|
if age_sec > max_age_sec:
|
|
return None
|
|
data = dict(entry["data"])
|
|
data["_age_ms"] = int(age_sec * 1000)
|
|
return data
|
|
|
|
def is_connected(self) -> bool:
|
|
return bool(self._connected and self._authenticated)
|
|
|
|
def subscribed_count(self) -> int:
|
|
with self._sub_lock:
|
|
return len(self._subscribed)
|
|
|
|
def attach_candle_aggregator(self, agg: Any) -> None:
|
|
"""KIS ``CandleAggregator`` 연결 — 키움 0B 틱으로 분봉 RAM 집계.
|
|
|
|
``set_candle_tick_codes()`` 로 후보 종목만 필터링하지 않으면
|
|
구독된 모든 종목 틱이 집계기로 들어감 (검증 전용 모드에서는 필터 권장).
|
|
"""
|
|
self._candle_agg = agg
|
|
logger.info("✅ KiwoomWebSocket → CandleAggregator 연결")
|
|
|
|
def set_candle_tick_codes(self, codes: Optional[Set[str]]) -> None:
|
|
"""집계기로 보낼 종목 코드. None=전체 구독 종목, 비어있지 않은 Set=해당 코드만."""
|
|
self._candle_agg_codes = codes
|
|
|
|
def attach_tick_recorder(self, recorder: Any) -> None:
|
|
"""TickRecorder 연결 — 키움 0B 체결 틱."""
|
|
self._tick_recorder = recorder
|
|
logger.info("✅ KiwoomWebSocket → TickRecorder 연결")
|
|
|
|
def attach_trigger_snapshot_recorder(self, recorder: Any) -> None:
|
|
"""TriggerSnapshotRecorder 연결 — 키움 0D/0w 스냅샷 DB."""
|
|
self._trigger_snapshot_recorder = recorder
|
|
logger.info("✅ KiwoomWebSocket → TriggerSnapshotRecorder 연결")
|
|
|
|
# ------------------------------------------------------------------
|
|
# 내부: 메인 수신 루프
|
|
# ------------------------------------------------------------------
|
|
def _run_loop(self) -> None:
|
|
while self._running:
|
|
try:
|
|
self._connect_and_serve()
|
|
except Exception as e:
|
|
logger.warning("키움 WS 루프 예외: %s", e)
|
|
finally:
|
|
self._connected = False
|
|
self._authenticated = False
|
|
if self._running:
|
|
self._reconnect_with_backoff()
|
|
|
|
def _connect_and_serve(self) -> None:
|
|
"""단일 연결 수명. blocking. 끊기면 반환 → 호출자가 backoff 후 재호출."""
|
|
if not self._ws_lib:
|
|
return
|
|
|
|
# 토큰 발급
|
|
token = self._get_kiwoom_token()
|
|
if not token:
|
|
logger.warning("⚠️ 키움 토큰 발급 실패 → WS 연결 보류 (60s)")
|
|
time.sleep(60)
|
|
return
|
|
|
|
# WebSocketApp 생성
|
|
self._ws = self._ws_lib.WebSocketApp(
|
|
self._ws_url,
|
|
on_open=self._on_open(token),
|
|
on_message=self._on_message,
|
|
on_error=self._on_error,
|
|
on_close=self._on_close,
|
|
)
|
|
self._last_connect_time = time.time()
|
|
# blocking — 연결 종료까지 여기서 대기
|
|
# ping_interval=0 : 키움은 JSON {"trnm":"PING"} keep-alive (프로토콜 ping 비호환)
|
|
self._ws.run_forever(ping_interval=0)
|
|
|
|
def _on_open(self, token: str):
|
|
"""on_open 콜백 팩토리 — token 캡처 후 LOGIN 발송."""
|
|
def _handler(ws):
|
|
self._connected = True
|
|
self._authenticated = False
|
|
try:
|
|
ws.send(json.dumps({"trnm": "LOGIN", "token": token}))
|
|
logger.info("📡 키움 WS 연결 → LOGIN 발송")
|
|
except Exception as e:
|
|
logger.warning("키움 WS LOGIN 발송 실패: %s", e)
|
|
return _handler
|
|
|
|
def _on_message(self, ws, message: str) -> None:
|
|
"""수신 메시지 디스패치 (LOGIN ack / REG ack / REAL / PING)."""
|
|
try:
|
|
msg = json.loads(message)
|
|
except Exception:
|
|
return
|
|
|
|
trnm = msg.get("trnm", "")
|
|
|
|
if trnm == "PING":
|
|
# 키움 PING → 그대로 echo (서버 정책)
|
|
try:
|
|
ws.send(message)
|
|
except Exception:
|
|
pass
|
|
return
|
|
|
|
if trnm == "LOGIN":
|
|
rc = msg.get("return_code")
|
|
rm = msg.get("return_msg", "")
|
|
if rc == 0:
|
|
self._authenticated = True
|
|
logger.info("✅ 키움 WS LOGIN OK")
|
|
# 디바운스 타이머 취소 — LOGIN 직후 일괄 REG 와 이중 전송 방지
|
|
with self._reg_timer_lock:
|
|
if self._reg_timer:
|
|
try:
|
|
self._reg_timer.cancel()
|
|
except Exception:
|
|
pass
|
|
self._reg_timer = None
|
|
# 누적된 구독 일괄 등록 (청크+간격 — TRNM=REG 레이트리밋)
|
|
with self._sub_lock:
|
|
pending = sorted(self._subscribed)
|
|
self._reg_batch_codes.clear()
|
|
if pending:
|
|
self._send_reg_chunked(pending)
|
|
# 조건검색 등 공유 세션 모듈 — LOGIN 직후 CNSRLST 재등록
|
|
self._fire_login_callbacks(ws)
|
|
else:
|
|
logger.warning("❌ 키움 WS LOGIN 실패 rc=%s msg=%s", rc, rm)
|
|
try:
|
|
ws.close()
|
|
except Exception:
|
|
pass
|
|
return
|
|
|
|
if trnm == "REG":
|
|
rc = msg.get("return_code")
|
|
if rc != 0:
|
|
logger.warning("⚠️ 키움 WS REG 실패: %s", msg.get("return_msg", ""))
|
|
return
|
|
|
|
if trnm == "REMOVE":
|
|
return # 무관심
|
|
|
|
if trnm == "REAL":
|
|
self._handle_real(msg)
|
|
self._dispatch_ext_handlers("REAL", ws, msg)
|
|
return
|
|
|
|
# 조건검색 CNSRLST / CNSRREQ / CNSRCLR 등
|
|
self._dispatch_ext_handlers(trnm, ws, msg)
|
|
|
|
def _handle_real(self, msg: dict) -> None:
|
|
"""실시간 데이터 처리 — 0B 체결 + 0D 호가잔량 + 0w 프로그램매매."""
|
|
items = msg.get("data") or []
|
|
for item in items:
|
|
sub_type = str(item.get("type", "")).strip()
|
|
code = str(item.get("item", "")).strip()
|
|
values = item.get("values") or {}
|
|
if sub_type == self.SUB_TYPE:
|
|
self._cache_tick(code, values)
|
|
elif sub_type == self.SUB_TYPE_ORDERBOOK and self._orderbook_cache:
|
|
try:
|
|
snap = self._orderbook_cache.update_from_kiwoom_0d(code, values)
|
|
if self._trigger_snapshot_recorder is not None:
|
|
tt = str(values.get("20") or values.get(self.FID_TICK_TIME) or "").strip()
|
|
self._trigger_snapshot_recorder.on_orderbook(snap, snap_time=tt or None)
|
|
except Exception as ex:
|
|
logger.debug("키움 0D 파싱 실패 %s: %s", code, ex)
|
|
elif sub_type == self.SUB_TYPE_PROGRAM and self._program_cache:
|
|
try:
|
|
snap = self._program_cache.update_from_kiwoom_0w(code, values)
|
|
if self._trigger_snapshot_recorder is not None:
|
|
tt = str(values.get("20") or "").strip()
|
|
self._trigger_snapshot_recorder.on_program(snap, snap_time=tt or None)
|
|
except Exception as ex:
|
|
logger.debug("키움 0w 파싱 실패 %s: %s", code, ex)
|
|
|
|
def _cache_tick(self, code: str, values: dict) -> None:
|
|
"""0B 체결 → 메모리 캐시 갱신 (KIS 와 동일 포맷)."""
|
|
if not code or not values:
|
|
return
|
|
try:
|
|
# 키움은 가격 부호로 등락 표시 — 절대값 취함
|
|
price_raw = values.get(self.FID_PRICE, "0")
|
|
price = abs(float(str(price_raw).replace(",", "")))
|
|
if price <= 0:
|
|
return
|
|
|
|
def _abs_str(v: str) -> str:
|
|
try:
|
|
return str(int(abs(float(str(v).replace(",", "")))))
|
|
except (ValueError, TypeError):
|
|
return "0"
|
|
|
|
# KIS inquire_price 호환 필드 (스칼라 dict)
|
|
cntr_raw = values.get(self.FID_EXEC_STRENGTH, "0")
|
|
try:
|
|
cntr_str = float(str(cntr_raw).replace(",", ""))
|
|
except (ValueError, TypeError):
|
|
cntr_str = 0.0
|
|
data_compat = {
|
|
"stck_prpr": str(int(price)),
|
|
"stck_oprc": _abs_str(values.get(self.FID_OPEN, "0")),
|
|
"stck_hgpr": _abs_str(values.get(self.FID_HIGH, "0")),
|
|
"stck_lwpr": _abs_str(values.get(self.FID_LOW, "0")),
|
|
"prdy_vrss": _abs_str(values.get(self.FID_CHANGE, "0")),
|
|
"prdy_ctrt": str(values.get(self.FID_CHANGE_PCT, "0")),
|
|
"cntr_str": str(cntr_str),
|
|
}
|
|
with self._cache_lock:
|
|
self._cache[code] = {"data": data_compat, "ts": time.time()}
|
|
|
|
# tick_time / tick_vol — CandleAggregator·TickRecorder 공용 (보유-only도 recorder 수집)
|
|
tt_raw = str(values.get(self.FID_TICK_TIME, "") or "").strip()
|
|
if len(tt_raw) >= 6:
|
|
tick_time = tt_raw[-6:]
|
|
else:
|
|
import datetime as _dt
|
|
tick_time = _dt.datetime.now().strftime("%H%M%S")
|
|
try:
|
|
tick_vol = int(
|
|
abs(float(str(values.get(self.FID_TICK_VOL, "0")).replace(",", "")))
|
|
)
|
|
except (ValueError, TypeError):
|
|
tick_vol = 0
|
|
|
|
# ── CandleAggregator: 후보만 (tick_to_agg = 후보−보유) ─────
|
|
if self._candle_agg is not None:
|
|
filt = self._candle_agg_codes
|
|
if filt is None or code in filt:
|
|
try:
|
|
self._candle_agg.on_tick(code, price, tick_vol, tick_time)
|
|
except Exception as ex:
|
|
logger.debug("키움→CandleAggregator on_tick 실패 %s: %s", code, ex)
|
|
|
|
# ── TickRecorder: 보유 포함 전 구독 종목 (백테 틱청산·B안 폴백) ─────
|
|
if self._tick_recorder is not None:
|
|
try:
|
|
self._tick_recorder.on_tick(
|
|
code, price, tick_vol, tick_time, source="kiwoom",
|
|
)
|
|
except Exception as ex:
|
|
logger.debug("키움→TickRecorder on_tick 실패 %s: %s", code, ex)
|
|
except (ValueError, TypeError) as e:
|
|
logger.debug("키움 0B 파싱 오류 %s: %s", code, e)
|
|
|
|
def _on_error(self, ws, error) -> None:
|
|
self._connected = False
|
|
self._authenticated = False
|
|
logger.warning("⚠️ 키움 WS 오류: %s", error)
|
|
|
|
def _on_close(self, ws, close_status_code, close_msg) -> None:
|
|
self._connected = False
|
|
self._authenticated = False
|
|
logger.info("🔌 키움 WS 연결 종료 (code=%s msg=%s)",
|
|
close_status_code, close_msg or "")
|
|
|
|
# ------------------------------------------------------------------
|
|
# 내부: 등록/해지 메시지 발송
|
|
# ------------------------------------------------------------------
|
|
def _send_reg_chunked(self, codes: List[str]) -> None:
|
|
"""REG 를 청크 단위로 전송하고 청크 사이에 sleep (키움 TRNM 건수 제한 회피)."""
|
|
if not codes:
|
|
return
|
|
chunk = self._reg_chunk_size()
|
|
gap = self._reg_gap_sec()
|
|
for i in range(0, len(codes), chunk):
|
|
part = codes[i:i + chunk]
|
|
self._send_reg(part)
|
|
if i + chunk < len(codes):
|
|
time.sleep(gap)
|
|
|
|
def _send_reg(self, codes: list) -> bool:
|
|
"""REG 발송 (그룹 1, 0B+0D 타입). 한 번에 여러 종목 OK."""
|
|
if not codes or not self._ws:
|
|
return False
|
|
try:
|
|
self._ws.send(json.dumps({
|
|
"trnm": "REG",
|
|
"grp_no": self.GROUP_NO,
|
|
"refresh": "1", # 재시작 시 등록 유지
|
|
"data": [{"item": list(codes), "type": self._reg_types()}],
|
|
}))
|
|
logger.info(
|
|
"📡 키움 WS REG 발송: %d종목 types=%s (총 %d/%d)",
|
|
len(codes), self._reg_types(), self.subscribed_count(),
|
|
self._max_subscriptions(),
|
|
)
|
|
return True
|
|
except Exception as e:
|
|
logger.warning("키움 WS REG 실패: %s", e)
|
|
return False
|
|
|
|
def _send_remove(self, codes: list) -> bool:
|
|
"""REMOVE 발송."""
|
|
if not codes or not self._ws:
|
|
return False
|
|
try:
|
|
self._ws.send(json.dumps({
|
|
"trnm": "REMOVE",
|
|
"grp_no": self.GROUP_NO,
|
|
"data": [{"item": list(codes), "type": self._reg_types()}],
|
|
}))
|
|
return True
|
|
except Exception as e:
|
|
logger.warning("키움 WS REMOVE 실패: %s", e)
|
|
return False
|
|
|
|
# ------------------------------------------------------------------
|
|
# 내부: 토큰
|
|
# ------------------------------------------------------------------
|
|
def _get_kiwoom_token(self) -> Optional[str]:
|
|
"""``kis_ws.KiwoomTokenManager`` 싱글톤 풀 재사용."""
|
|
try:
|
|
from .kis_ws import _get_kiwoom_token_cached
|
|
except ImportError:
|
|
logger.warning("kis_ws._get_kiwoom_token_cached import 실패")
|
|
return None
|
|
return _get_kiwoom_token_cached(self.app_key, self.app_secret, self.is_mock)
|
|
|
|
# ------------------------------------------------------------------
|
|
# 내부: 재연결 백오프 (KIS WS 와 동일 정책)
|
|
# ------------------------------------------------------------------
|
|
def _reconnect_with_backoff(self) -> None:
|
|
now = time.time()
|
|
|
|
# 5분 이상 안정 연결 후 끊긴 거면 카운터 초기화 (정상 재연결 케이스)
|
|
if (self._last_connect_time > 0
|
|
and now - self._last_connect_time >= self.STABLE_CONN_RESET_SEC):
|
|
self._reconnect_count = 0
|
|
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
|
|
|
|
# 시간당 한도 초과 → 1시간 대기
|
|
cutoff = now - 3600
|
|
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])
|
|
logger.warning("⚠️ 키움 WS 시간당 재연결 한도 초과 → %ds 대기", int(wait))
|
|
time.sleep(max(60.0, wait))
|
|
return
|
|
|
|
# 총 한도 초과 → 비활성
|
|
if self._reconnect_count >= self.MAX_RECONNECT_ATTEMPTS:
|
|
logger.warning("⛔ 키움 WS 총 재연결 한도 초과 → 자동 비활성")
|
|
self._running = False
|
|
return
|
|
|
|
delay = min(self._reconnect_delay, self.RECONNECT_MAX_DELAY_SEC)
|
|
logger.info("⏳ 키움 WS %ds 후 재연결 시도 (#%d)",
|
|
int(delay), self._reconnect_count + 1)
|
|
time.sleep(delay)
|
|
self._reconnect_times.append(time.time())
|
|
self._reconnect_count += 1
|
|
self._reconnect_delay = min(self._reconnect_delay * 2, self.RECONNECT_MAX_DELAY_SEC)
|