Changes: - Added new API endpoints for continuing and confirming Optuna jobs, allowing for better management of ongoing studies. - Introduced detailed logging for tick feed tracking and order book processing, improving traceability of vendor performance during backtests. - Updated database schema to include new fields for managing Optuna study results, enhancing the ability to track study progress and outcomes. - Refactored existing functions to utilize the new logging and tracking features, ensuring consistency across the backtesting framework. Impact: - These enhancements improve the robustness and transparency of the Optuna backtesting process, facilitating better analysis and optimization of trading strategies.
1119 lines
48 KiB
Python
1119 lines
48 KiB
Python
"""
|
|
kis_trader/ws/kiwoom_ws.py — 키움 WebSocket 실시간 시세 캐시
|
|
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
|
|
|
|
목적
|
|
----
|
|
실매 키움 실시간 시세(체결 0B · 호가 0D · 프로그램 0w) **한 소켓**.
|
|
예전에는 KIS 41 한도 대비 검증용이었으나, 지금은 틱/호가 본체.
|
|
|
|
``LIVE_TICK_PROVIDER`` 는 이 모듈이 아니라 ``WSManager.get_price`` 읽기 순서.
|
|
``WS_SUBSCRIBE_KIS_MINIMAL`` 은 한투 구독 명단. 키움 후보 구독과 무관.
|
|
|
|
설계 원칙 (KIS WS 와 동일)
|
|
--------------------------
|
|
- ``Dict[str, Dict]`` 메모리 캐시 (락만 잠그고 마이크로초 read/write)
|
|
- ``get_price(code)`` 인터페이스를 KIS WS 와 100% 동일 포맷으로 제공
|
|
- ``KiwoomTokenManager`` 싱글톤 (kis_ws.py 안에 있음) 재사용
|
|
- 재연결 백오프, 종목 등록/해지
|
|
|
|
키움 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)
|
|
|
|
토글 (DB env_config)
|
|
--------------------
|
|
키움 WS 기동: Orchestrator ``_start_ws_validator`` — 키 있으면 기동.
|
|
``WS_PROVIDER=kis_only`` 로 끄지 않음(구식 주석).
|
|
``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) 초당 허용 건수 초과 방지(배치·전송 간격·단건 디바운스).
|
|
``KIWOOM_WS_REG_REFRESH`` (기본 1) — REG ``refresh``. 0=기존 그룹 구독 해지 후 재등록.
|
|
LOGIN 직후 일괄 REG **첫 청크만** 0. 장중 종목 추가 REG는 항상 1
|
|
(청크마다 0이면 앞 종목이 지워짐). 장중 상시 0 비추(catch-up 재점화).
|
|
"""
|
|
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 = "228" # 체결강도 (%) — 키움 0B 공식
|
|
FID_UPPER_LIMIT_TIME = "567" # 상한가발생시간 HHmmss
|
|
FID_LOWER_LIMIT_TIME = "568" # 하한가발생시간 HHmmss
|
|
|
|
# 한도 (키움 권장)
|
|
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()
|
|
# 종목별 마지막 FID20 — 재연결 catch-up(시간 후퇴/지연) vs 시계 동결 구분
|
|
self._fid20_last: Dict[str, Any] = {}
|
|
# 현재 FID20 를 처음 본 벽시계 — 같은 초 연속 체결 vs 아침 동결(수 분) 구분
|
|
self._fid20_first_wall: Dict[str, float] = {}
|
|
self._stale_skip_log_ts: Dict[str, float] = {}
|
|
self._price_listeners: list = []
|
|
self._price_listener_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 _unsubscribe_all_before_close(self, reason: str = "", *, clear_ram: bool = True) -> int:
|
|
"""세션 close 전 서버 REMOVE. LOGIN 전(구독 없음)이면 0.
|
|
|
|
RAM 은 ``clear_ram=True`` (stop) 일 때만 비운다.
|
|
"""
|
|
with self._sub_lock:
|
|
codes = sorted(self._subscribed)
|
|
if clear_ram:
|
|
self._reg_batch_codes.clear()
|
|
self._subscribed.clear()
|
|
if not codes or not self._connected or not self._authenticated or self._ws is None:
|
|
return 0
|
|
n_ok = 0
|
|
try:
|
|
chunk = self._reg_chunk_size()
|
|
gap = self._reg_gap_sec()
|
|
for i in range(0, len(codes), chunk):
|
|
part = codes[i:i + chunk]
|
|
if self._send_remove(part):
|
|
n_ok += len(part)
|
|
if i + chunk < len(codes):
|
|
time.sleep(gap)
|
|
logger.info(
|
|
"✅ 키움 WS 종료 전 REMOVE %d/%d종목 (%s)",
|
|
n_ok, len(codes), reason or "-",
|
|
)
|
|
except Exception as e:
|
|
logger.warning("키움 WS 종료 전 REMOVE 실패: %s", e)
|
|
return n_ok
|
|
|
|
def stop(self) -> None:
|
|
"""수신 스레드 종료 + 구독 REMOVE 후 소켓 닫기 (재시작 세션 꼬임 완화)."""
|
|
with self._reg_timer_lock:
|
|
if self._reg_timer:
|
|
try:
|
|
self._reg_timer.cancel()
|
|
except Exception:
|
|
pass
|
|
self._reg_timer = None
|
|
try:
|
|
self._unsubscribe_all_before_close(reason="stop", clear_ram=True)
|
|
except Exception as e:
|
|
logger.warning("키움 WS 종료 전 REMOVE 실패: %s", e)
|
|
self._running = False
|
|
self._connected = False
|
|
self._authenticated = False
|
|
try:
|
|
if self._ws is not None:
|
|
self._ws.close()
|
|
except Exception:
|
|
pass
|
|
logger.info("⏹ 키움 WS 종료")
|
|
|
|
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:
|
|
# 사용자의 혼동 방지를 위해 KIWOOM_WS_ORDERBOOK_ENABLED 대신,
|
|
# LIVE_OB_PROVIDER가 kiwoom이거나, 키움 호가 적재(WS_ORDERBOOK_SAVE_KIWOOM)가 켜져있으면 자동 수신.
|
|
from ..utils.env import get_env_from_db, get_env_bool
|
|
live_ob_provider = (get_env_from_db("LIVE_OB_PROVIDER", "kiwoom") or "kiwoom").strip().lower()
|
|
save_kiwoom = get_env_bool("WS_ORDERBOOK_SAVE_KIWOOM", True)
|
|
if live_ob_provider == "kiwoom" or save_kiwoom:
|
|
return True
|
|
return get_env_bool("KIWOOM_WS_ORDERBOOK_ENABLED", False)
|
|
|
|
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 add_price_listener(self, callback) -> None:
|
|
"""현재가 갱신 콜백 등록. callback(code, price, data_dict)."""
|
|
if callback is None:
|
|
return
|
|
with self._price_listener_lock:
|
|
if callback not in self._price_listeners:
|
|
self._price_listeners.append(callback)
|
|
|
|
def remove_price_listener(self, callback) -> None:
|
|
if callback is None:
|
|
return
|
|
with self._price_listener_lock:
|
|
try:
|
|
self._price_listeners.remove(callback)
|
|
except ValueError:
|
|
pass
|
|
|
|
def _emit_price_listeners(self, code: str, price: float, data: Dict) -> None:
|
|
with self._price_listener_lock:
|
|
listeners = list(self._price_listeners)
|
|
for cb in listeners:
|
|
try:
|
|
cb(code, price, data)
|
|
except Exception as ex:
|
|
logger.debug("price_listener 예외 %s: %s", code, ex)
|
|
|
|
def get_price(self, code: str, max_age_sec: Optional[float] = None) -> Optional[Dict]:
|
|
"""KIS WS ``get_price`` 와 동일 포맷 반환. None = 마지막 RAM (체결 공백이어도 유지).
|
|
|
|
반환::
|
|
|
|
{
|
|
"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 max_age_sec is not None and 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)")
|
|
try:
|
|
from kis_trader.utils.ops_alert import ops_alert
|
|
ops_alert(
|
|
"token_kiwoom",
|
|
"키움 토큰 발급 실패",
|
|
detail="시세 WS·조건검색 불가 → 60s 후 재시도",
|
|
level="critical",
|
|
session_only=False,
|
|
)
|
|
except Exception:
|
|
pass
|
|
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
|
|
# 안정 LOGIN 성공 → 재연결 카운터 리셋 (주말 Bye 후 소모분 복구)
|
|
self._reconnect_count = 0
|
|
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
|
|
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:
|
|
# LOGIN ack 를 받은 같은 on_message 스레드에서 REG 청크 전송.
|
|
# 이 동안 0B 파싱이 밀림. 장중 추가는 Timer 디바운스(별 스레드).
|
|
self._send_reg_chunked(pending, group_replace=True)
|
|
# 조건검색 등 공유 세션 모듈 — LOGIN 직후 CNSRLST 재등록
|
|
self._fire_login_callbacks(ws)
|
|
else:
|
|
logger.warning("❌ 키움 WS LOGIN 실패 rc=%s msg=%s", rc, rm)
|
|
try:
|
|
from kis_trader.utils.ops_alert import ops_alert
|
|
ops_alert(
|
|
"ws_kiwoom_login",
|
|
"키움 WS LOGIN 실패",
|
|
detail=f"rc={rc} msg={rm}",
|
|
level="critical",
|
|
session_only=False,
|
|
)
|
|
except Exception:
|
|
pass
|
|
# 805004 등 서버 토큰 거부 → 캐시 만료 전이라도 무효화 후 재발급
|
|
if self._is_token_rejected_login(rc, rm):
|
|
self._invalidate_token_cache(
|
|
"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 호환 + 키움 0B 원본 FID 보존 (분석·추후 필터용)
|
|
exec_raw = values.get(self.FID_EXEC_STRENGTH, "0")
|
|
try:
|
|
cntr_str = float(str(exec_raw).replace(",", ""))
|
|
except (ValueError, TypeError):
|
|
cntr_str = 0.0
|
|
upper_limit_raw = str(values.get(self.FID_UPPER_LIMIT_TIME, "") or "").strip()
|
|
lower_limit_raw = str(values.get(self.FID_LOWER_LIMIT_TIME, "") or "").strip()
|
|
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),
|
|
"upper_limit_time": upper_limit_raw,
|
|
"lower_limit_time": lower_limit_raw,
|
|
"kiwoom_fid228": str(exec_raw),
|
|
"kiwoom_fid567": upper_limit_raw,
|
|
"kiwoom_fid568": lower_limit_raw,
|
|
"chetime": "",
|
|
"kiwoom_fid20": "",
|
|
}
|
|
|
|
# tick_time / tick_vol — CandleAggregator·TickRecorder 공용 (보유-only도 recorder 수집)
|
|
# 늦음 판정 = 이 틱의 체결시각(FID20) vs 우리 서버 지금. 증권사끼리 비교 아님.
|
|
# 매매 RAM/리스너/봉: 읽기나이(LIVE_FEED_FALLBACK, 기본 3초) 넘으면 안 넣음 → 2차·3차가 연다.
|
|
# 적재(ws_ticks): 버리지 않고 전부. LIVE_MAX<=0 이면 적재 skip 없음.
|
|
# 아침 FID20 동결: 같은 FID20 가 TIME_MAX초 이상 안 움직일 때만 RAM 유지.
|
|
import datetime as _dt
|
|
tt_raw = str(values.get(self.FID_TICK_TIME, "") or "").strip()
|
|
_now_dt = _dt.datetime.now()
|
|
_now_wall = time.time()
|
|
if len(tt_raw) >= 6:
|
|
tick_time_pkt = tt_raw[-6:]
|
|
else:
|
|
tick_time_pkt = _now_dt.strftime("%H%M%S")
|
|
data_compat["chetime"] = tick_time_pkt
|
|
data_compat["kiwoom_fid20"] = tick_time_pkt
|
|
data_compat["tick_time"] = (
|
|
tt_raw[:14] if len(tt_raw) >= 14 else (_now_dt.strftime("%Y%m%d") + tick_time_pkt)
|
|
)
|
|
_pkt_dt = None
|
|
try:
|
|
if len(tt_raw) >= 14:
|
|
_pkt_dt = _dt.datetime.strptime(tt_raw[:14], "%Y%m%d%H%M%S")
|
|
elif len(tt_raw) >= 6:
|
|
_pkt_dt = _dt.datetime.strptime(
|
|
_now_dt.strftime("%Y%m%d") + tick_time_pkt, "%Y%m%d%H%M%S",
|
|
)
|
|
except Exception:
|
|
_pkt_dt = None
|
|
_lag_sec = 0.0
|
|
if _pkt_dt is not None:
|
|
_lag_sec = float((_now_dt - _pkt_dt).total_seconds())
|
|
try:
|
|
_live_max = int(get_env_int("KIWOOM_TICK_LIVE_MAX_LAG_SEC", 0) or 0)
|
|
except (TypeError, ValueError):
|
|
_live_max = 0
|
|
try:
|
|
from kis_trader.engine.feed_fallback import live_feed_fallback_max_age_sec
|
|
_read_max = float(live_feed_fallback_max_age_sec())
|
|
except Exception:
|
|
_read_max = 2.0
|
|
_time_max = max(5, int(get_env_int("KIWOOM_TICK_TIME_MAX_LAG_SEC", 120)))
|
|
_last_fid = self._fid20_last.get(code)
|
|
if _pkt_dt is not None and _last_fid != _pkt_dt:
|
|
self._fid20_first_wall[code] = _now_wall
|
|
_first_w = float(self._fid20_first_wall.get(code, _now_wall) or _now_wall)
|
|
# 같은 FID20 가 TIME_MAX초 이상 지속 = 아침 시계 멈춤. 밀린 줄의 같은 초 2건(수십 ms)은 제외.
|
|
_frozen = bool(
|
|
_pkt_dt is not None
|
|
and _last_fid is not None
|
|
and _pkt_dt == _last_fid
|
|
and (_now_wall - _first_w) >= float(_time_max)
|
|
)
|
|
# 재연결 버퍼: FID20가 지금보다 읽기나이(2초) 이상 과거면 매매 RAM 미반영.
|
|
# 적재는 기본 전부(LIVE_MAX=0). 양수일 때만 그 초 초과 미저장.
|
|
_skip_ram = bool(
|
|
_read_max > 0
|
|
and _pkt_dt is not None
|
|
and _lag_sec > float(_read_max)
|
|
and not _frozen
|
|
)
|
|
_skip_persist = bool(
|
|
_live_max > 0
|
|
and _pkt_dt is not None
|
|
and _lag_sec > float(_live_max)
|
|
and not _frozen
|
|
)
|
|
if _pkt_dt is not None:
|
|
self._fid20_last[code] = _pkt_dt
|
|
|
|
if not _skip_ram:
|
|
wall_ts = time.time()
|
|
with self._cache_lock:
|
|
self._cache[code] = {"data": data_compat, "ts": wall_ts}
|
|
self._emit_price_listeners(code, float(price), data_compat)
|
|
else:
|
|
_log_gap = max(5.0, float(get_env_float("KIWOOM_TICK_STALE_LOG_GAP_SEC", 30.0) or 30.0))
|
|
_prev_log = float(self._stale_skip_log_ts.get(code, 0.0) or 0.0)
|
|
if (time.time() - _prev_log) >= _log_gap:
|
|
self._stale_skip_log_ts[code] = time.time()
|
|
logger.warning(
|
|
"키움 0B stale skip RAM %s fid20=%s lag=%.0fs price=%.0f (매수체크 미반영·2차 폴백)",
|
|
code, tick_time_pkt, _lag_sec, price,
|
|
)
|
|
|
|
tick_time = tick_time_pkt
|
|
if _pkt_dt is None:
|
|
tick_time = _now_dt.strftime("%H%M%S")
|
|
elif _frozen and _lag_sec > float(_time_max):
|
|
# FID20 동결(아침 55분 버그) — 분봉·틱청산 타임라인만 wall-clock
|
|
tick_time = _now_dt.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 = 후보−보유) ─────
|
|
# 밀린 메인 틱은 확정봉에 넣지 않음 (freeze 오염 방지)
|
|
if self._candle_agg is not None and not _skip_ram:
|
|
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, source="kiwoom")
|
|
except Exception as ex:
|
|
logger.debug("키움→CandleAggregator on_tick 실패 %s: %s", code, ex)
|
|
|
|
# ── TickRecorder: 보유 포함 전 구독 종목 (백테 틱청산·B안 폴백) ─────
|
|
# 공책은 밀려도 적재. 읽기 2초와 분리. 옵투나는 lag>2초 메인을 보조로 교체.
|
|
if self._tick_recorder is not None and not _skip_persist:
|
|
try:
|
|
self._tick_recorder.on_tick(
|
|
code, price, tick_vol, tick_time, source="kiwoom",
|
|
cntr_str=cntr_str,
|
|
upper_limit_time=upper_limit_raw or None,
|
|
tick_time_raw=tt_raw or None,
|
|
)
|
|
except Exception as ex:
|
|
logger.debug("키움→TickRecorder on_tick 실패 %s: %s", code, ex)
|
|
|
|
# ── 호가 틱동기: 체결 1건당 RAM 0D 스냅 1장 (LS tick 모드와 동일) ──
|
|
# 밀린 체결시각에 지금 호가를 붙이면 옵투나 호가가 미래가 됨 → RAM과 같이 skip.
|
|
if self._trigger_snapshot_recorder is not None and not _skip_ram:
|
|
try:
|
|
max_age = float(
|
|
getattr(
|
|
self._trigger_snapshot_recorder,
|
|
"orderbook_tick_max_age_sec",
|
|
3.0,
|
|
)
|
|
or 3.0
|
|
)
|
|
snap = self.get_orderbook_snapshot(code, max_age_sec=max_age)
|
|
if snap is not None:
|
|
self._trigger_snapshot_recorder.on_orderbook_tick_sync(
|
|
code, snap, snap_time=tick_time,
|
|
)
|
|
except Exception as ex:
|
|
logger.debug("키움→호가틱동기 실패 %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 _reg_refresh_keep(self) -> str:
|
|
"""REG refresh. 기본 1(기존 유지). 0=그룹 기존 item/type 해지 — 테스트용."""
|
|
v = int(get_env_int("KIWOOM_WS_REG_REFRESH", 1) or 1)
|
|
return "0" if v == 0 else "1"
|
|
|
|
def _send_reg_chunked(self, codes: List[str], *, group_replace: bool = False) -> None:
|
|
"""REG 를 청크 단위로 전송하고 청크 사이에 sleep (키움 TRNM 건수 제한 회피).
|
|
|
|
group_replace=True (LOGIN 일괄): ``KIWOOM_WS_REG_REFRESH=0`` 이면 **첫 청크만** refresh=0.
|
|
이후 청크·장중 추가 REG는 항상 1 (0을 반복하면 앞 청크 구독이 삭제됨).
|
|
"""
|
|
if not codes:
|
|
return
|
|
chunk = self._reg_chunk_size()
|
|
gap = self._reg_gap_sec()
|
|
want = self._reg_refresh_keep() if group_replace else "1"
|
|
if group_replace and want == "0":
|
|
logger.warning(
|
|
"⚠️ 키움 WS REG refresh=0 테스트 — LOGIN 일괄 첫 청크만 기존 구독 해지 후 재등록"
|
|
)
|
|
for i in range(0, len(codes), chunk):
|
|
part = codes[i:i + chunk]
|
|
refresh = want if i == 0 else "1"
|
|
self._send_reg(part, refresh=refresh)
|
|
if i + chunk < len(codes):
|
|
time.sleep(gap)
|
|
|
|
def _send_reg(self, codes: list, *, refresh: str = "1") -> bool:
|
|
"""REG 발송 (그룹 1, 0B+0D 타입). 한 번에 여러 종목 OK."""
|
|
if not codes or not self._ws or not self._connected:
|
|
return False
|
|
flag = "0" if str(refresh).strip() == "0" else "1"
|
|
try:
|
|
self._ws.send(json.dumps({
|
|
"trnm": "REG",
|
|
"grp_no": self.GROUP_NO,
|
|
"refresh": flag,
|
|
"data": [{"item": list(codes), "type": self._reg_types()}],
|
|
}))
|
|
logger.info(
|
|
"📡 키움 WS REG 발송: %d종목 types=%s refresh=%s (총 %d/%d)",
|
|
len(codes), self._reg_types(), flag, 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 or not self._connected:
|
|
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)
|
|
|
|
@staticmethod
|
|
def _is_token_rejected_login(rc: Any, rm: Any) -> bool:
|
|
"""WS LOGIN 응답이 '토큰 무효'인지 — 캐시 invalidate 대상."""
|
|
try:
|
|
if int(rc) == 805004:
|
|
return True
|
|
except Exception:
|
|
pass
|
|
m = str(rm or "")
|
|
return ("Token이 유효하지 않습니다" in m) or ("토큰이 유효하지 않습니다" in m)
|
|
|
|
def _invalidate_token_cache(self, reason: str) -> None:
|
|
try:
|
|
from .kis_ws import _invalidate_kiwoom_token_cached
|
|
except ImportError:
|
|
logger.warning("kis_ws._invalidate_kiwoom_token_cached import 실패")
|
|
return
|
|
_invalidate_kiwoom_token_cached(
|
|
self.app_key, self.app_secret, self.is_mock, reason=reason,
|
|
)
|
|
|
|
# ------------------------------------------------------------------
|
|
# 내부: 재연결 백오프 (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
|
|
|
|
# 총 한도 초과 — 기본은 긴 쿨다운 후 카운터 리셋 (주말 토큰거부 연타로
|
|
# 월요일까지 영구 OFF 되면 SCALP/MOMENTUM 조건식 REAL 이 끊김).
|
|
# 응급으로만 KIWOOM_WS_HARD_DISABLE_ON_MAX_RECONNECT=true.
|
|
if self._reconnect_count >= self.MAX_RECONNECT_ATTEMPTS:
|
|
if get_env_bool("KIWOOM_WS_HARD_DISABLE_ON_MAX_RECONNECT", False):
|
|
logger.warning("⛔ 키움 WS 총 재연결 한도 초과 → 자동 비활성")
|
|
self._running = False
|
|
return
|
|
cool = float(get_env_float("KIWOOM_WS_MAX_RECONNECT_COOLDOWN_SEC", 1800.0))
|
|
cool = max(60.0, min(cool, 86400.0))
|
|
logger.warning(
|
|
"⚠️ 키움 WS 총 재연결 한도(#%d) → %.0fs 쿨다운 후 카운터 리셋 "
|
|
"(영구 OFF 안 함 · HARD_DISABLE=true 시만 비활성)",
|
|
self._reconnect_count, cool,
|
|
)
|
|
# 쿨다운 전 토큰 캐시 비우기 — 죽은 토큰 재전송 루프 차단
|
|
self._invalidate_token_cache("max reconnect cooldown")
|
|
time.sleep(cool)
|
|
self._reconnect_count = 0
|
|
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
|
|
self._reconnect_times.clear()
|
|
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)
|