Files
kis_bot/kis_trader/ws/kiwoom_ws.py
Your Name 1d69c217e2 fix(정합성): 틱 lag wall-clock 정합 + 3벤더 DB 저장 스위치 통일
- feed_fallback.bar_is_garbage: 봉끝 기준 → 각 틱의 recv_ts wall-clock 기준으로 정정
  실매 RAM 3초컷과 동일 논리 → 유동성 낮은 종목 부당 스킵 해소
- candle_garbage_fallback_enabled: 기본 True 복원 (wall-clock 정정 후 안전)
- param_search_optuna·run_tail_backtest_cli: CANDLE_GARBAGE_FALLBACK·BACKTEST_USE_RUST
  강제 os.environ 세팅 제거 → DB env·CLI 플래그로만 관리 (UI 존중)
- WS_TICK_DB_SAVE_LAG_CUT_ENABLED 신설 (bool, 기본 false, 3벤더 공통)
  OFF=키움/KIS/LS 모든 틱 lag 무관 전부 저장 (벤더 통계·재현·백테 정합)
  ON=lag > LIVE_FEED_FALLBACK_MAX_AGE_SEC 이면 미저장 (미래 A안)
- KIWOOM_TICK_LIVE_MAX_LAG_SEC 완전 폐기 → 위 스위치로 통일
- kiwoom_ws: _skip_persist 로직 새 스위치로 교체
- kis_ws·ls_ws: _skip_persist_kis/_ls 신규 (벤더별 상이했던 정책 통일)
- docs/정합성.md §9 신설 (문제·결정·시나리오·향후 A안 전환법)
- docs/rust_engine_parity_port_plan.md (신규 설계)

실매 스모크: logs/test_live_execution_validation_20260906_191845.log
  → 최종: 통과 · 👑 완결
브라우저 검증: http://192.168.0.149:5050/#liveconfig → 새 스위치 노출, JS 오류 없음

영향: 실매(DB 저장 정책 통일, RAM 컷 변경 없음) + 백테/Optuna(실매 정합)

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-06 19:23:01 +09:00

1159 lines
50 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")
from .ws_reconnect_backoff import ws_reconnect_delay_for_attempt
# ── 환경변수 헬퍼 (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]
from .kiwoom_ws_diag import log_kiwoom_api_msg, log_kiwoom_ws
# ──────────────────────────────────────────────────────────────────────
# 메인 클래스
# ──────────────────────────────────────────────────────────────────────
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" # 종목프로그램매매
# 재연결 sleep — ws_reconnect_backoff (기본 1,3,5,7,10초)
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._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:
log_kiwoom_ws(
logger, "send_skip",
level="warning",
trnm=str(msg.get("trnm") or ""),
extra="미인증 또는 ws=None",
force=True,
)
return False
try:
self._ws.send(json.dumps(msg))
return True
except Exception as e:
log_kiwoom_ws(
logger, "send_fail",
level="warning",
trnm=str(msg.get("trnm") or ""),
extra=str(e),
force=True,
)
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 in ("CNSRLST", "CNSRREQ", "CNSRCLR", "REG", "REMOVE", "LOGIN"):
log_kiwoom_api_msg(logger, str(trnm or "?"), msg)
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
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 실패 rc=%s msg=%s",
rc, 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 snap is not None and 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())
# 2026-09-06: KIWOOM_TICK_LIVE_MAX_LAG_SEC 삭제 → WS_TICK_DB_SAVE_LAG_CUT_ENABLED 통일 (3벤더 공통).
# 스위치 OFF(기본)=전부 저장(벤더 통계 유지). ON=lag > LIVE_FEED_FALLBACK 이면 미저장.
try:
_db_cut_on = bool(get_env_bool("WS_TICK_DB_SAVE_LAG_CUT_ENABLED", False))
except Exception:
_db_cut_on = False
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 미반영.
_skip_ram = bool(
_read_max > 0
and _pkt_dt is not None
and _lag_sec > float(_read_max)
and not _frozen
)
# DB 저장: 스위치 OFF(기본)면 전부 저장. ON 이면 RAM 컷과 동일 기준(_read_max) 사용.
_skip_persist = bool(
_db_cut_on
and _pkt_dt is not None
and _lag_sec > float(_read_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
err_s = str(error or "")
log_kiwoom_ws(
logger, "on_error",
level="warning",
extra=err_s,
force=True,
)
logger.warning("⚠️ 키움 WS 오류: %s", error)
def _on_close(self, ws, close_status_code, close_msg) -> None:
self._connected = False
self._authenticated = False
log_kiwoom_ws(
logger, "on_close",
level="info",
rc=close_status_code,
msg=close_msg or "",
force=True,
)
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
# 시간당 한도 초과 → 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])
log_kiwoom_ws(
logger, "reconnect_hourly_cap",
level="warning",
extra=f"wait={int(wait)}s count={len(self._reconnect_times)}",
force=True,
)
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_times.clear()
return
attempt = self._reconnect_count + 1
delay = ws_reconnect_delay_for_attempt(attempt)
logger.info("⏳ 키움 WS %ds 후 재연결 시도 (#%d)",
int(delay), attempt)
time.sleep(delay)
self._reconnect_times.append(time.time())
self._reconnect_count += 1