- 백테스트 및 실거래 시 틱과 호가의 벤더 출처(ob_source, entry_source) 기록 및 추적 강화 (tail_engine.py) - 웹 UI '체결디버그'에 [틱:kis / 호가:ls] 형태로 데이터 출처를 직관적으로 표출 (backtest.js, backtest.html) - LS WebSocket 구독 100건 제한 하드코딩 해제 및 env_config_ext 연동 (ls_ws.py) - 기타 백테스트 웹 및 DB 관련 최적화 적용
1905 lines
76 KiB
Python
1905 lines
76 KiB
Python
"""
|
|
kis_trader/ws/ls_ws.py — LS증권 WebSocket 시세 캐시 (그림자 검증용)
|
|
==================================================================
|
|
|
|
목적
|
|
----
|
|
실매 시세(KIS/키움)와 **병렬**로 LS 실키 WS 를 붙여
|
|
가격·틱·1분봉을 **별도 테이블**에 쌓고 갭을 비교한다.
|
|
**매매 의사결정에는 기본 사용하지 않는다** (읽기·적재 전용).
|
|
``LS_WS_BLOCK_BUY_WHILE_RECOVERING=true`` 일 때만 복구 중 매수 게이트에 참여.
|
|
|
|
안정성 (B/A/C/D)
|
|
----------------
|
|
- B: 모든 WS send 는 ``_send_lock`` + **timeout** 직렬화. REG 는 워커 큐.
|
|
- A: ``LS_WS_TR_MODE=us3``(기본) → 통합 ``US3`` + ``U``+패딩.
|
|
``s3k3`` → ``S3_``/``K3_`` (실험용).
|
|
- C: 생존 = **프로토콜 ping**(기본 20s). 앱 JSON 하트비트는 보내지 않음(스펙 없음).
|
|
틱 silence 강제재연결 = **벽시계 우선** (국장 정규창 / 해외 hold).
|
|
sticky JIF ``jstatus=21`` 만으로 장외 워치독 금지.
|
|
``JIF`` 마감 상태는 벽시계 안에서만 추가 차단.
|
|
- D: 재연결·REG 복구 중 ``is_recovering()`` — 매수 게이트는 env 로 선택.
|
|
- E: 국내/해외 **hold 창**(``ls_ws_session_windows``) — 갭·장외는 소켓 close 후 대기.
|
|
세션 전환 시 **oauth 재발급 금지**(공용 토큰 무효화 방지, 키움 과거 사고와 동일 유형).
|
|
|
|
토글: ``LS_WS_VALIDATION_ENABLED`` (기본 false)
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import queue
|
|
import threading
|
|
import time
|
|
from collections import defaultdict, deque
|
|
from datetime import datetime
|
|
from typing import Any, Callable, Dict, List, Optional, Set, Tuple
|
|
|
|
import requests
|
|
|
|
logger = logging.getLogger("LSWebSocket")
|
|
|
|
from .ws_reconnect_backoff import ws_reconnect_delay_for_attempt
|
|
|
|
try:
|
|
from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int
|
|
except ImportError:
|
|
def get_env_bool(key, default=False): # type: ignore[misc]
|
|
return default
|
|
|
|
def get_env_int(key, default=0): # type: ignore[misc]
|
|
return int(default)
|
|
|
|
def get_env_float(key, default=0.0): # type: ignore[misc]
|
|
return float(default)
|
|
|
|
def get_env_from_db(key, default=""): # type: ignore[misc]
|
|
return default
|
|
|
|
LS_REST_BASE = "https://openapi.ls-sec.co.kr:8080"
|
|
LS_WS_REAL = "wss://openapi.ls-sec.co.kr:9443/websocket"
|
|
LS_WS_MOCK = "wss://openapi.ls-sec.co.kr:29443/websocket"
|
|
|
|
# JIF jstatus — 정규장 체결이 기대되는 상태 (스펙: 21=장시작)
|
|
_JIF_STATUS_EXPECT_TICKS = frozenset({"21"})
|
|
|
|
# AFR 세션 sticky — 장중 다른 jstatus(66/70 등)가 덮어써도 개장 sticky 유지.
|
|
# 실측(2026-07-31): 08:50~ 25/24/23/22 → 09:00 21 / 마감 30·31·41~44
|
|
_JIF_AFR_OPEN_DEFAULT = frozenset({"21", "22", "23", "24", "25"})
|
|
_JIF_AFR_CLOSE_DEFAULT = frozenset({"30", "31", "41", "42", "43", "44"})
|
|
|
|
|
|
def _parse_jif_status_set(env_key: str, default: frozenset) -> frozenset:
|
|
raw = str(get_env_from_db(env_key, "") or "").strip()
|
|
if not raw:
|
|
return default
|
|
out = {p.strip() for p in raw.replace(";", ",").split(",") if p.strip()}
|
|
return frozenset(out) if out else default
|
|
|
|
TickRecorder = Callable[[str, Dict[str, Any]], None]
|
|
CandleFlusher = Callable[[str, Dict[str, Any]], None]
|
|
OrderbookRecorder = Callable[[str, Dict[str, Any]], None]
|
|
ViRecorder = Callable[[str, Dict[str, Any]], None]
|
|
|
|
# 매수 게이트·진단용 — start() 시 등록, stop() 시 해제
|
|
_active_ls_ws: Optional["LSWebSocketPriceCache"] = None
|
|
_active_lock = threading.Lock()
|
|
|
|
|
|
def get_active_ls_ws() -> Optional["LSWebSocketPriceCache"]:
|
|
with _active_lock:
|
|
return _active_ls_ws
|
|
|
|
|
|
def domestic_unified_tr_key(shcode: str, width: int = 10) -> str:
|
|
code = (shcode or "").strip()
|
|
raw = f"U{code}" if (len(code) == 6 and code.isdigit()) else (
|
|
code if code.startswith("U") else f"U{code}"
|
|
)
|
|
return raw if len(raw) >= width else raw.ljust(width)
|
|
|
|
|
|
def overseas_tr_key(exchcd: str, symbol: str, width: int = 18) -> str:
|
|
raw = f"{exchcd}{symbol}"
|
|
return raw if len(raw) >= width else raw.ljust(width)
|
|
|
|
|
|
def fetch_ls_access_token(
|
|
app_key: str,
|
|
app_secret: str,
|
|
timeout: float = 15.0,
|
|
*,
|
|
force: bool = False,
|
|
reason: str = "",
|
|
) -> str:
|
|
"""``POST /oauth2/token`` — expires_in 기반 공유 캐시 (한도 준수 재사용)."""
|
|
from kis_trader.network.ls_token import fetch_ls_access_token as _shared
|
|
|
|
return _shared(
|
|
app_key, app_secret, timeout=timeout, force=force, reason=reason or "ls_ws",
|
|
)
|
|
|
|
|
|
class LSWebSocketPriceCache:
|
|
"""LS WS 가격 캐시 — get_price 포맷을 KIS/키움과 맞춤 (stck_prpr, _age_ms)."""
|
|
|
|
def __init__(
|
|
self,
|
|
app_key: str,
|
|
app_secret: str,
|
|
*,
|
|
is_mock: bool = False,
|
|
also_hoga: bool = False,
|
|
) -> None:
|
|
try:
|
|
import websocket # websocket-client
|
|
except ImportError as e:
|
|
raise RuntimeError("websocket-client 필요") from e
|
|
self._websocket = websocket
|
|
self.app_key = app_key
|
|
self.app_secret = app_secret
|
|
self.is_mock = bool(is_mock)
|
|
self.also_hoga = bool(also_hoga)
|
|
self.ws_url = LS_WS_MOCK if self.is_mock else LS_WS_REAL
|
|
|
|
self._token: str = ""
|
|
self._token_at: float = 0.0
|
|
self._ws: Any = None
|
|
self._thread: Optional[threading.Thread] = None
|
|
self._reg_thread: Optional[threading.Thread] = None
|
|
self._watch_thread: Optional[threading.Thread] = None
|
|
self._session_guard_thread: Optional[threading.Thread] = None
|
|
self._running = False
|
|
self._opened = threading.Event()
|
|
|
|
self._sub_lock = threading.Lock()
|
|
self._send_lock = threading.Lock() # SSL write 직렬화 (BAD_LENGTH 레이스 방지)
|
|
self._force_closing = False # UNREG→close 재진입 방지
|
|
self._hold_ram_extra_codes_fn: Optional[Any] = None # WSManager → 보유 pin
|
|
self._ws_manager_bound: bool = False
|
|
self._subscribed: Set[str] = set() # 6자리 KR 코드 (소켓 실구독)
|
|
# code → {"permanent","condition","default",...} — 해제는 owner 전부일 때만
|
|
self._sub_owners: Dict[str, Set[str]] = {}
|
|
self._us_subscribed: Set[str] = set() # 티커
|
|
self._us_sub_owners: Dict[str, Set[str]] = {}
|
|
self._cache_lock = threading.Lock()
|
|
self._cache: Dict[str, Dict[str, Any]] = {}
|
|
self._market_cache: Dict[str, str] = {} # code → K|Q|E
|
|
|
|
self._reg_q: queue.Queue = queue.Queue()
|
|
self._recovering = False
|
|
self._recovering_lock = threading.Lock()
|
|
self._last_tick_mono: float = 0.0
|
|
self._opened_mono: float = 0.0
|
|
self._watchdog_empty_streak: int = 0
|
|
self._ws_reconnect_step: int = 0
|
|
# Bye/조기 CLOSE 서킷 — OPEN 직후 끊김 시 백오프 리셋 금지·토큰 1회 갱신
|
|
self._early_close_streak: int = 0
|
|
self._last_close_lived_sec: float = 0.0
|
|
self._last_close_early: bool = False
|
|
self._pending_bye: bool = False
|
|
self._stable_open: bool = False
|
|
|
|
# JIF 장운영정보 (앱 하트비트 아님 — 세션 판정용)
|
|
self._jif_lock = threading.Lock()
|
|
self._jangubun: str = ""
|
|
self._jstatus: str = ""
|
|
self._jif_mono: float = 0.0
|
|
# 국장 AFR용 sticky (해외 밤장과 무관 — LS JIF 만)
|
|
self._kr_afr_session_live: bool = False
|
|
self._kr_afr_session_src: str = ""
|
|
self._kr_afr_opened_ymd: str = "" # 당일 개장전/장시작 JIF 수신일
|
|
self._kr_afr_closed_ymd: str = "" # 당일 마감 JIF 수신일
|
|
|
|
# 1분봉 롤업 (code → bucket) + 확정봉 RAM (전략 get_candles 호환)
|
|
self._candle_lock = threading.Lock()
|
|
self._candles: Dict[str, Dict[str, Any]] = {}
|
|
# (code, tf_min) → 확정봉 deque (candle_time 스키마)
|
|
self._confirmed: Dict[Tuple[str, int], deque] = defaultdict(
|
|
lambda: deque(maxlen=max(50, get_env_int("LS_WS_CONFIRMED_MAX", 500)))
|
|
)
|
|
self._tick_recorder: Optional[TickRecorder] = None
|
|
self._candle_flusher: Optional[CandleFlusher] = None
|
|
self._orderbook_recorder: Optional[OrderbookRecorder] = None
|
|
self._vi_recorder: Optional[ViRecorder] = None
|
|
self._candle_save_filter: Optional[Callable[[str], bool]] = None
|
|
self._orderbook_cache: Any = None
|
|
self._ob_last_save_mono: Dict[str, float] = {}
|
|
# BaseStrategy 틱매도 등 — KIS/키움과 동일 시그니처
|
|
self._price_listeners: List[Callable] = []
|
|
# code → last VI-active monotonic (해제·만료 시 discard)
|
|
self._vi_lock = threading.Lock()
|
|
self._vi_active_mono: Dict[str, float] = {}
|
|
try:
|
|
from kis_trader.ws.orderbook_cache import OrderbookCache
|
|
|
|
self._orderbook_cache = OrderbookCache()
|
|
except Exception as e:
|
|
logger.debug("LS OrderbookCache 미사용: %s", e)
|
|
|
|
# ── public API (키움/KIS 와 동일 취지) ──────────────────────────────
|
|
def attach_tick_recorder(self, fn: TickRecorder) -> None:
|
|
self._tick_recorder = fn
|
|
|
|
def attach_candle_flusher(self, fn: CandleFlusher) -> None:
|
|
self._candle_flusher = fn
|
|
|
|
def attach_orderbook_recorder(self, fn: OrderbookRecorder) -> None:
|
|
"""LS UH1/H1_/HA_ → DB 등 (실매 호가필터 경로와 분리)."""
|
|
self._orderbook_recorder = fn
|
|
|
|
def attach_vi_recorder(self, fn: ViRecorder) -> None:
|
|
"""LS UVI/VI_ 발동·해제 → ls_ws_vi 등."""
|
|
self._vi_recorder = fn
|
|
|
|
def set_candle_save_filter(self, fn: Optional[Callable[[str], bool]]) -> None:
|
|
"""분봉 롤업·flush 대상. None=전원. 후보는 영구만(후보 분봉 폭주 방지)."""
|
|
self._candle_save_filter = fn
|
|
|
|
def is_in_vi(self, code: str) -> bool:
|
|
"""종목이 VI 발동 중인지 (해제 누락 대비 stale 만료)."""
|
|
c = self._normalize_kr_code(code)
|
|
if not c:
|
|
return False
|
|
stale = max(30, int(get_env_int("LS_WS_VI_STALE_SEC", 180) or 180))
|
|
now = time.monotonic()
|
|
with self._vi_lock:
|
|
mono = self._vi_active_mono.get(c)
|
|
if mono is None:
|
|
return False
|
|
if now - mono > stale:
|
|
self._vi_active_mono.pop(c, None)
|
|
return False
|
|
return True
|
|
|
|
def get_orderbook_snapshot(self, code: str, max_age_sec: float = 3.0):
|
|
"""키움 ``get_orderbook_snapshot`` 호환 — LS RAM 호가."""
|
|
if self._orderbook_cache is None:
|
|
return None
|
|
return self._orderbook_cache.get(code, max_age_sec=max_age_sec)
|
|
|
|
def is_recovering(self) -> bool:
|
|
with self._recovering_lock:
|
|
return bool(self._recovering)
|
|
|
|
def blocks_new_buy(self) -> bool:
|
|
"""복구 중 신규매수 차단 — env ON 일 때만 True."""
|
|
if not get_env_bool("LS_WS_BLOCK_BUY_WHILE_RECOVERING", False):
|
|
return False
|
|
return self.is_recovering() or not self.is_connected()
|
|
|
|
def start(self) -> bool:
|
|
global _active_ls_ws
|
|
if self._thread and self._thread.is_alive():
|
|
return True
|
|
try:
|
|
self._token = fetch_ls_access_token(
|
|
self.app_key, self.app_secret, force=False, reason="ls_ws_start",
|
|
)
|
|
self._token_at = time.time()
|
|
except Exception as e:
|
|
logger.error("LS 토큰 발급 실패: %s", e)
|
|
return False
|
|
self._running = True
|
|
self._opened.clear()
|
|
self._last_tick_mono = time.monotonic()
|
|
self._set_recovering(True)
|
|
self._thread = threading.Thread(
|
|
target=self._run_forever, name="LS-WS", daemon=True
|
|
)
|
|
self._reg_thread = threading.Thread(
|
|
target=self._reg_worker, name="LS-WS-REG", daemon=True
|
|
)
|
|
self._watch_thread = threading.Thread(
|
|
target=self._watchdog_loop, name="LS-WS-Watch", daemon=True
|
|
)
|
|
self._session_guard_thread = threading.Thread(
|
|
target=self._session_guard_loop, name="LS-WS-SessionGuard", daemon=True
|
|
)
|
|
self._thread.start()
|
|
self._reg_thread.start()
|
|
self._watch_thread.start()
|
|
self._session_guard_thread.start()
|
|
with _active_lock:
|
|
_active_ls_ws = self
|
|
# hold 외 기동: OPEN 대기하지 않음 — 루프가 다음 hold 까지 슬립 (재발급 없음)
|
|
if not self._should_hold_socket():
|
|
from kis_trader.utils.ls_ws_session_windows import (
|
|
seconds_until_ls_socket_open,
|
|
)
|
|
|
|
wait_sec = seconds_until_ls_socket_open(
|
|
n_us_subscribed=self._n_us_subscribed(),
|
|
)
|
|
logger.info(
|
|
"✅ LS WS 기동(hold 외 대기) — 다음 hold까지 약 %.0f분 "
|
|
"(소켓 미연결·토큰 캐시 유지)",
|
|
wait_sec / 60.0,
|
|
)
|
|
return True
|
|
if not self._opened.wait(timeout=15.0):
|
|
logger.error("LS WS OPEN timeout")
|
|
self.stop()
|
|
return False
|
|
logger.info(
|
|
"✅ LS WS 연결 (%s mock=%s tr_mode=%s ping=%ss)",
|
|
self.ws_url,
|
|
self.is_mock,
|
|
self._tr_mode(),
|
|
int(get_env_int("LS_WS_PING_INTERVAL_SEC", 20) or 0),
|
|
)
|
|
return True
|
|
|
|
def _graceful_unreg_all(
|
|
self,
|
|
*,
|
|
clear_ram: bool = True,
|
|
ram_keep_codes: Optional[Set[str]] = None,
|
|
) -> int:
|
|
"""종료 직전 서버 구독 UNREG — 재시작 시 세션 꼬임/한도 거부 완화.
|
|
|
|
REG 워커 큐가 아니라 동기 전송(갭 준수). 타임아웃 초과 시 남은 종목은
|
|
TCP close 에 맡긴다 (systemd stop 지연·폭주 방지).
|
|
``clear_ram=True`` + ``ram_keep_codes``: 서버 UNREG 후 RAM 은 keep 만 유지
|
|
(hold reopen 시 07:00 mass REG 방지 — permanent·보유 pin).
|
|
"""
|
|
if self._ws is None or not self._opened.is_set():
|
|
return 0
|
|
with self._sub_lock:
|
|
kr = list(self._subscribed)
|
|
us = list(self._us_subscribed)
|
|
if not kr and not us:
|
|
# JIF 만 구독 중일 수 있음
|
|
pass
|
|
timeout = max(
|
|
1.0,
|
|
float(get_env_float("LS_WS_STOP_UNREG_TIMEOUT_SEC", 8.0) or 8.0),
|
|
)
|
|
gap_ms = max(
|
|
10,
|
|
int(
|
|
get_env_int(
|
|
"LS_WS_STOP_UNREG_GAP_MS",
|
|
int(get_env_int("LS_WS_REG_GAP_MS", 80) or 80),
|
|
)
|
|
or 30
|
|
),
|
|
)
|
|
deadline = time.monotonic() + timeout
|
|
n_ok = 0
|
|
timed_out = False
|
|
|
|
def _one(tr_type: str, tr_cd: str, tr_key: str) -> bool:
|
|
nonlocal n_ok, timed_out
|
|
if time.monotonic() >= deadline:
|
|
timed_out = True
|
|
return False
|
|
if not self._opened.is_set() or self._ws is None:
|
|
return False
|
|
self._send_typed(tr_type, tr_cd, tr_key, close_on_fail=False)
|
|
n_ok += 1
|
|
time.sleep(gap_ms / 1000.0)
|
|
return True
|
|
|
|
# 장운영 JIF 해지
|
|
if get_env_bool("LS_WS_JIF_ENABLED", True):
|
|
jif_key = (get_env_from_db("LS_WS_JIF_TR_KEY", "0") or "0").strip() or "0"
|
|
_one("4", "JIF", jif_key)
|
|
|
|
for code in kr:
|
|
if timed_out:
|
|
break
|
|
tr_cd, tr_key, hoga_cd, hoga_key = self._kr_tr_pair(code)
|
|
if not _one("4", tr_cd, tr_key):
|
|
break
|
|
if self.also_hoga and not _one("4", hoga_cd, hoga_key):
|
|
break
|
|
if get_env_bool("LS_WS_UVI_ENABLED", True):
|
|
vi_cd, vi_key = self._vi_tr_pair(code)
|
|
if not _one("4", vi_cd, vi_key):
|
|
break
|
|
|
|
for sym in us:
|
|
if timed_out:
|
|
break
|
|
if not _one("4", "GSC", overseas_tr_key("82", sym)):
|
|
break
|
|
|
|
if clear_ram:
|
|
keep = {
|
|
str(c).strip()
|
|
for c in (ram_keep_codes or set())
|
|
if str(c).strip()
|
|
}
|
|
with self._sub_lock:
|
|
if keep:
|
|
for c in list(self._subscribed):
|
|
if c not in keep:
|
|
self._subscribed.discard(c)
|
|
self._sub_owners.pop(c, None)
|
|
for c in list(self._us_subscribed):
|
|
if c not in keep:
|
|
self._us_subscribed.discard(c)
|
|
self._us_sub_owners.pop(c, None)
|
|
else:
|
|
self._subscribed.clear()
|
|
self._sub_owners.clear()
|
|
self._us_subscribed.clear()
|
|
self._us_sub_owners.clear()
|
|
|
|
if timed_out:
|
|
logger.warning(
|
|
"⚠️ LS WS 종료 UNREG 타임아웃 %.1fs — sends=%d KR=%d US=%d "
|
|
"(잔여 서버구독은 close 에 위임)",
|
|
timeout, n_ok, len(kr), len(us),
|
|
)
|
|
elif n_ok > 0:
|
|
logger.info(
|
|
"✅ LS WS 종료 전 UNREG 완료 sends=%d KR=%d US=%d gap=%dms",
|
|
n_ok, len(kr), len(us), gap_ms,
|
|
)
|
|
return n_ok
|
|
|
|
def stop(self) -> None:
|
|
"""구독 UNREG → 소켓 close — 봇 재시작 시 서버 세션 잔존 완화."""
|
|
global _active_ls_ws
|
|
try:
|
|
self._graceful_unreg_all()
|
|
except Exception as e:
|
|
logger.warning("LS WS 종료 전 UNREG 실패: %s", e)
|
|
self._running = False
|
|
self._set_recovering(False)
|
|
try:
|
|
self._reg_q.put_nowait(None) # poison
|
|
except Exception:
|
|
pass
|
|
if self._ws is not None:
|
|
try:
|
|
self._ws.close()
|
|
except Exception:
|
|
pass
|
|
for th in (
|
|
self._reg_thread,
|
|
self._watch_thread,
|
|
self._session_guard_thread,
|
|
self._thread,
|
|
):
|
|
if th is not None and th.is_alive():
|
|
try:
|
|
th.join(timeout=2.0)
|
|
except Exception:
|
|
pass
|
|
with _active_lock:
|
|
if _active_ls_ws is self:
|
|
_active_ls_ws = None
|
|
logger.info("⏹ LS WS 종료")
|
|
|
|
def is_connected(self) -> bool:
|
|
return self._opened.is_set() and self._running
|
|
|
|
def subscribe(self, code: str, owner: str = "default") -> bool:
|
|
"""시세 구독. ``owner`` 가 남아 있으면 중복 REG 안 함.
|
|
|
|
owner 예: ``permanent``(영구구독), ``condition``(LS 조건 이력), ``default``,
|
|
``spill``(타벤더 한도/실패 즉시 폴백 — RAM 전용).
|
|
|
|
성공·이미구독 True / 빈코드·미연결·로컬한도 False (매니저 spill 체인용).
|
|
"""
|
|
code = (code or "").strip()
|
|
if not code:
|
|
return False
|
|
own = (owner or "default").strip() or "default"
|
|
# spill 즉시성: 미연결이면 대기 없이 False (일반 owner 는 기존처럼 집합 적재 후 REG 큐)
|
|
if own == "spill" and not self.is_connected():
|
|
return False
|
|
max_n = int(get_env_int("LS_WS_MAX_SUBSCRIPTIONS", 100) or 0)
|
|
if code.isdigit() and len(code) == 6:
|
|
need_reg = False
|
|
with self._sub_lock:
|
|
ow = self._sub_owners.setdefault(code, set())
|
|
if own in ow and code in self._subscribed:
|
|
return True
|
|
ow.add(own)
|
|
if code not in self._subscribed:
|
|
self._subscribed.add(code)
|
|
need_reg = True
|
|
if need_reg:
|
|
self._enqueue_kr_reg(code)
|
|
return True
|
|
else:
|
|
need_reg = False
|
|
with self._sub_lock:
|
|
ow = self._us_sub_owners.setdefault(code, set())
|
|
if own in ow and code in self._us_subscribed:
|
|
return True
|
|
ow.add(own)
|
|
if code not in self._us_subscribed:
|
|
self._us_subscribed.add(code)
|
|
need_reg = True
|
|
if need_reg:
|
|
self._enqueue_send("3", "GSC", overseas_tr_key("82", code))
|
|
return True
|
|
|
|
def unsubscribe(self, code: str, owner: str = "default") -> None:
|
|
"""owner 제거. 다른 owner 가 남으면 소켓 구독 유지."""
|
|
code = (code or "").strip()
|
|
if not code:
|
|
return
|
|
own = (owner or "default").strip() or "default"
|
|
if code.isdigit() and len(code) == 6:
|
|
drop_socket = False
|
|
with self._sub_lock:
|
|
ow = self._sub_owners.get(code)
|
|
if ow is not None:
|
|
ow.discard(own)
|
|
if ow:
|
|
return
|
|
self._sub_owners.pop(code, None)
|
|
if code not in self._subscribed:
|
|
return
|
|
self._subscribed.discard(code)
|
|
drop_socket = True
|
|
if not drop_socket:
|
|
return
|
|
tr_cd, tr_key, hoga_cd, hoga_key = self._kr_tr_pair(code)
|
|
self._enqueue_send("4", tr_cd, tr_key)
|
|
if self.also_hoga:
|
|
self._enqueue_send("4", hoga_cd, hoga_key)
|
|
if get_env_bool("LS_WS_UVI_ENABLED", True):
|
|
vi_cd, vi_key = self._vi_tr_pair(code)
|
|
self._enqueue_send("4", vi_cd, vi_key)
|
|
if self._orderbook_cache is not None:
|
|
self._orderbook_cache.remove(code)
|
|
with self._vi_lock:
|
|
self._vi_active_mono.pop(code, None)
|
|
else:
|
|
drop_socket = False
|
|
with self._sub_lock:
|
|
ow = self._us_sub_owners.get(code)
|
|
if ow is not None:
|
|
ow.discard(own)
|
|
if ow:
|
|
return
|
|
self._us_sub_owners.pop(code, None)
|
|
if code not in self._us_subscribed:
|
|
return
|
|
self._us_subscribed.discard(code)
|
|
drop_socket = True
|
|
if drop_socket:
|
|
self._enqueue_send("4", "GSC", overseas_tr_key("82", code))
|
|
|
|
def sync_owner_codes(self, owner: str, codes: Set[str]) -> None:
|
|
"""특정 owner 구독 집합을 ``codes`` 로 맞춤 (diff subscribe/unsubscribe)."""
|
|
own = (owner or "default").strip() or "default"
|
|
desired = {(c or "").strip() for c in (codes or set()) if (c or "").strip()}
|
|
with self._sub_lock:
|
|
cur: Set[str] = set()
|
|
for c, ow in self._sub_owners.items():
|
|
if own in ow:
|
|
cur.add(c)
|
|
for c, ow in self._us_sub_owners.items():
|
|
if own in ow:
|
|
cur.add(c)
|
|
add = sorted(desired - cur)
|
|
drop = sorted(cur - desired)
|
|
if add or drop or own == "condition":
|
|
logger.info(
|
|
"LS WS owner=%s sync +%d -%d desired=%d (socket_KR≈확인은 OPEN/워치독)",
|
|
own, len(add), len(drop), len(desired),
|
|
)
|
|
for c in add:
|
|
self.subscribe(c, owner=own)
|
|
for c in drop:
|
|
self.unsubscribe(c, owner=own)
|
|
|
|
def get_price(self, code: str, max_age_sec: Optional[float] = None) -> Optional[Dict[str, Any]]:
|
|
with self._cache_lock:
|
|
d = self._cache.get(code)
|
|
if not d:
|
|
return None
|
|
age = time.time() - float(d.get("_ts", 0))
|
|
if max_age_sec is not None and age > max_age_sec:
|
|
return None
|
|
out = dict(d)
|
|
out["_age_ms"] = int(age * 1000)
|
|
return out
|
|
|
|
def add_price_listener(self, callback) -> None:
|
|
if callback is None:
|
|
return
|
|
if callback not in self._price_listeners:
|
|
self._price_listeners.append(callback)
|
|
|
|
def remove_price_listener(self, callback) -> None:
|
|
if callback is None:
|
|
return
|
|
try:
|
|
self._price_listeners.remove(callback)
|
|
except ValueError:
|
|
pass
|
|
|
|
def _notify_price_listeners(self, code: str, price: float, data: Dict[str, Any]) -> None:
|
|
for cb in list(self._price_listeners):
|
|
try:
|
|
cb(code, price, data)
|
|
except Exception:
|
|
pass
|
|
|
|
@staticmethod
|
|
def _forming_to_strategy_bar(cur: Dict[str, Any]) -> Dict[str, Any]:
|
|
"""LS forming/flush dict → CandleAggregator 호환."""
|
|
from kis_trader.network.ls_chart import ls_datetime_to_candle_time
|
|
|
|
dt_s = str(cur.get("datetime") or "")
|
|
ct = ls_datetime_to_candle_time(dt_s)
|
|
return {
|
|
"candle_time": ct,
|
|
"open": float(cur.get("open") or 0),
|
|
"high": float(cur.get("high") or 0),
|
|
"low": float(cur.get("low") or 0),
|
|
"close": float(cur.get("close") or 0),
|
|
"volume": float(cur.get("volume") or 0),
|
|
"tick_count": int(cur.get("tick_count") or 0),
|
|
"is_confirmed": 0,
|
|
"source": "ls",
|
|
"rsi_2": None,
|
|
"rsi_3": None,
|
|
"rsi_5": None,
|
|
}
|
|
|
|
@staticmethod
|
|
def _calc_rsi(closes: list, period: int) -> Optional[float]:
|
|
if len(closes) < period + 1:
|
|
return None
|
|
recent = closes[-(period + 1):]
|
|
gains, losses = [], []
|
|
for i in range(1, len(recent)):
|
|
delta = recent[i] - recent[i - 1]
|
|
gains.append(max(delta, 0))
|
|
losses.append(max(-delta, 0))
|
|
avg_gain = sum(gains) / period if period > 0 else 0
|
|
avg_loss = sum(losses) / period if period > 0 else 0
|
|
if avg_loss == 0:
|
|
return 100.0
|
|
rs = avg_gain / avg_loss
|
|
return round(100 - (100 / (1 + rs)), 2)
|
|
|
|
def _attach_rsi(self, bars: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
|
closes = [float(b.get("close") or 0) for b in bars]
|
|
out = []
|
|
for i, b in enumerate(bars):
|
|
row = dict(b)
|
|
sub = closes[: i + 1]
|
|
row["rsi_2"] = self._calc_rsi(sub, 2)
|
|
row["rsi_3"] = self._calc_rsi(sub, 3)
|
|
row["rsi_5"] = self._calc_rsi(sub, 5)
|
|
out.append(row)
|
|
return out
|
|
|
|
def _push_confirmed(self, code: str, bar: Dict[str, Any]) -> None:
|
|
"""확정봉 RAM 적재 (동일 candle_time 이면 갱신)."""
|
|
tf = max(1, int(bar.get("tf_min") or get_env_int("LS_WS_CANDLE_TF_MIN", 1)))
|
|
ct = str(bar.get("candle_time") or "")[:12]
|
|
if len(ct) < 12:
|
|
return
|
|
key = (code, tf)
|
|
with self._candle_lock:
|
|
buf = self._confirmed[key]
|
|
if buf and str(buf[-1].get("candle_time") or "")[:12] == ct:
|
|
buf[-1] = dict(bar)
|
|
buf[-1]["is_confirmed"] = 1
|
|
else:
|
|
row = dict(bar)
|
|
row["is_confirmed"] = 1
|
|
buf.append(row)
|
|
|
|
def get_candles(self, code: str, tf: int, n: int = 50) -> list:
|
|
"""최근 n개 확정 봉 (오래된→최신). RAM → DB 폴백."""
|
|
tf_i = max(1, int(tf or 1))
|
|
lim = max(1, int(n or 50))
|
|
with self._candle_lock:
|
|
buf = list(self._confirmed.get((code, tf_i), []))
|
|
if len(buf) >= lim:
|
|
return self._attach_rsi(buf[-lim:])
|
|
# DB 폴백으로 보강
|
|
try:
|
|
from database import TradeDB
|
|
|
|
db_bars = TradeDB().get_ls_ws_candles(code, tf_min=tf_i, limit=lim)
|
|
except Exception:
|
|
db_bars = []
|
|
if not db_bars and not buf:
|
|
return []
|
|
# RAM 우선 merge by candle_time
|
|
by_ct: Dict[str, Dict[str, Any]] = {}
|
|
for b in db_bars:
|
|
ct = str(b.get("candle_time") or "")[:12]
|
|
if ct:
|
|
by_ct[ct] = dict(b)
|
|
for b in buf:
|
|
ct = str(b.get("candle_time") or "")[:12]
|
|
if ct:
|
|
by_ct[ct] = dict(b)
|
|
merged = [by_ct[k] for k in sorted(by_ct.keys())]
|
|
return self._attach_rsi(merged[-lim:])
|
|
|
|
def get_current_candle(self, code: str, tf: int) -> Optional[dict]:
|
|
"""진행 중 봉 (is_confirmed=0)."""
|
|
tf_i = max(1, int(tf or 1))
|
|
want_tf = max(1, get_env_int("LS_WS_CANDLE_TF_MIN", 1))
|
|
if tf_i != want_tf:
|
|
# 1분 롤업만 유지 — 다른 TF 요청은 None (키움 롤업과 동일 제약)
|
|
return None
|
|
with self._candle_lock:
|
|
cur = self._candles.get(code)
|
|
if not cur:
|
|
return None
|
|
return self._forming_to_strategy_bar(cur)
|
|
|
|
def fill_gap_from_rest(
|
|
self,
|
|
code: str,
|
|
*,
|
|
qrycnt: Optional[int] = None,
|
|
upsert_db: bool = True,
|
|
) -> int:
|
|
"""
|
|
t8412 분봉 → RAM 확정봉 + ls_ws_candles upsert.
|
|
반환: 반영한 봉 수. 실패/비활성 시 0.
|
|
"""
|
|
from kis_trader.network.ls_chart import (
|
|
candle_time_to_ls_datetime,
|
|
fetch_ls_minute_chart_df,
|
|
)
|
|
|
|
code = (code or "").strip()
|
|
if not code:
|
|
return 0
|
|
df = fetch_ls_minute_chart_df(code, ncnt=1, qrycnt=qrycnt)
|
|
if df is None or df.empty:
|
|
return 0
|
|
tf = max(1, get_env_int("LS_WS_CANDLE_TF_MIN", 1))
|
|
n_ok = 0
|
|
db = None
|
|
if upsert_db:
|
|
try:
|
|
from database import TradeDB
|
|
|
|
db = TradeDB()
|
|
except Exception as e:
|
|
logger.debug("LS gap DB 연결 실패: %s", e)
|
|
db = None
|
|
for _, row in df.iterrows():
|
|
ct = str(row.get("time") or "")[:12]
|
|
if len(ct) < 12:
|
|
continue
|
|
bar = {
|
|
"candle_time": ct,
|
|
"open": float(row["open"]),
|
|
"high": float(row["high"]),
|
|
"low": float(row["low"]),
|
|
"close": float(row["close"]),
|
|
"volume": float(row.get("volume") or 0),
|
|
"tick_count": 0,
|
|
"is_confirmed": 1,
|
|
"source": "ls_t8412",
|
|
"tf_min": tf,
|
|
}
|
|
self._push_confirmed(code, bar)
|
|
if db is not None:
|
|
try:
|
|
db.upsert_ls_ws_candle(
|
|
code=code,
|
|
candle={
|
|
"datetime": candle_time_to_ls_datetime(ct),
|
|
"tf_min": tf,
|
|
"open": bar["open"],
|
|
"high": bar["high"],
|
|
"low": bar["low"],
|
|
"close": bar["close"],
|
|
"volume": bar["volume"],
|
|
"tick_count": 0,
|
|
},
|
|
)
|
|
except Exception as e:
|
|
logger.debug("LS gap upsert %s %s: %s", code, ct, e)
|
|
n_ok += 1
|
|
if n_ok:
|
|
logger.info("📥 [LS갭] %s t8412 → %d봉 (RAM+DB)", code, n_ok)
|
|
return n_ok
|
|
|
|
# ── TR / 시장 ─────────────────────────────────────────────────────
|
|
def _tr_mode(self) -> str:
|
|
# 기본 us3: 당일 실측 전 틱이 US3. s3k3 전환 후 틱 0+watchdog 폭주 확인됨.
|
|
v = (get_env_from_db("LS_WS_TR_MODE", "us3") or "us3").strip().lower()
|
|
if v in ("s3k3", "s3", "k3", "split"):
|
|
return "s3k3"
|
|
return "us3"
|
|
|
|
def _lookup_market(self, code: str) -> str:
|
|
"""stock_meta.market: K=KOSPI, Q=KOSDAQ, E=ETF. 없으면 ''."""
|
|
if code in self._market_cache:
|
|
return self._market_cache[code]
|
|
m = ""
|
|
try:
|
|
from database import TradeDB
|
|
db = TradeDB()
|
|
row = db.conn.execute(
|
|
"SELECT market FROM stock_meta WHERE code=%s LIMIT 1",
|
|
(code,),
|
|
).fetchone()
|
|
if row:
|
|
m = str(row.get("market") or "").strip().upper()[:1]
|
|
except Exception as e:
|
|
logger.debug("LS stock_meta 조회 실패 %s: %s", code, e)
|
|
self._market_cache[code] = m
|
|
return m
|
|
|
|
def _kr_tr_pair(self, code: str) -> Tuple[str, str, str, str]:
|
|
"""(체결 tr_cd, tr_key, 호가 tr_cd, 호가 tr_key)."""
|
|
if self._tr_mode() == "us3":
|
|
k = domestic_unified_tr_key(code)
|
|
return "US3", k, "UH1", k
|
|
m = self._lookup_market(code)
|
|
if m == "K" or m == "E":
|
|
return "S3_", code, "H1_", code
|
|
if m == "Q":
|
|
return "K3_", code, "HA_", code
|
|
# 메타 없음: 통합 US3 폴백 (잘못된 S3_/K3_ 보다 안전)
|
|
k = domestic_unified_tr_key(code)
|
|
return "US3", k, "UH1", k
|
|
|
|
def _vi_tr_pair(self, code: str) -> Tuple[str, str]:
|
|
"""(VI tr_cd, tr_key) — us3→UVI+U패딩, s3k3→VI_+6자리."""
|
|
if self._tr_mode() == "us3":
|
|
return "UVI", domestic_unified_tr_key(code)
|
|
return "VI_", code
|
|
|
|
def _enqueue_kr_reg(self, code: str) -> None:
|
|
tr_cd, tr_key, hoga_cd, hoga_key = self._kr_tr_pair(code)
|
|
self._enqueue_send("3", tr_cd, tr_key)
|
|
if self.also_hoga:
|
|
self._enqueue_send("3", hoga_cd, hoga_key)
|
|
if get_env_bool("LS_WS_UVI_ENABLED", True):
|
|
vi_cd, vi_key = self._vi_tr_pair(code)
|
|
self._enqueue_send("3", vi_cd, vi_key)
|
|
|
|
def _enqueue_send(self, tr_type: str, tr_cd: str, tr_key: str) -> None:
|
|
try:
|
|
self._reg_q.put_nowait(("send", tr_type, tr_cd, tr_key))
|
|
except Exception as e:
|
|
logger.debug("LS REG 큐 put 실패: %s", e)
|
|
|
|
def _enqueue_replay(self) -> None:
|
|
try:
|
|
self._reg_q.put_nowait(("replay", None, None, None))
|
|
except Exception as e:
|
|
logger.debug("LS REG replay 큐 실패: %s", e)
|
|
|
|
def _set_recovering(self, on: bool) -> None:
|
|
with self._recovering_lock:
|
|
self._recovering = bool(on)
|
|
|
|
# ── internals ─────────────────────────────────────────────────────
|
|
def _ensure_token(self) -> str:
|
|
# 스펙 expires_in 공유 캐시 — force 금지(세션 전환·재연결에서 공용 토큰 무효화 방지)
|
|
self._token = fetch_ls_access_token(
|
|
self.app_key, self.app_secret, force=False, reason="ls_ws_ensure",
|
|
)
|
|
self._token_at = time.time()
|
|
return self._token
|
|
|
|
def _stable_open_sec(self) -> float:
|
|
return max(5.0, float(get_env_float("LS_WS_STABLE_OPEN_SEC", 30.0) or 30.0))
|
|
|
|
def _mark_stable_open(self, why: str) -> None:
|
|
"""REG 완료·장시간 OPEN 후에만 재연결 step/Bye streak 리셋."""
|
|
if self._stable_open and self._ws_reconnect_step == 0 and self._early_close_streak == 0:
|
|
return
|
|
self._stable_open = True
|
|
self._ws_reconnect_step = 0
|
|
self._early_close_streak = 0
|
|
self._pending_bye = False
|
|
logger.info(
|
|
"LS WS 안정 OPEN 확정 — reconnect/Bye streak 리셋 (%s)",
|
|
why,
|
|
)
|
|
|
|
def _refresh_token_after_bye(self) -> None:
|
|
"""조기 Bye 서킷: 로컬 캐시 무효화 후 1회 재발급 (min_gap 존중)."""
|
|
from kis_trader.network.ls_token import invalidate_ls_access_token
|
|
|
|
invalidate_ls_access_token(
|
|
self.app_key, self.app_secret, reason="ls_ws_early_bye",
|
|
)
|
|
self._token = ""
|
|
self._token = fetch_ls_access_token(
|
|
self.app_key,
|
|
self.app_secret,
|
|
force=True,
|
|
reason="ls_ws_early_bye",
|
|
)
|
|
self._token_at = time.time()
|
|
|
|
def _hold_ram_keep_codes(self) -> Set[str]:
|
|
"""hold 양보 후 RAM 유지 집합 — _permanent owner + (옵션) 보유."""
|
|
keep: Set[str] = set()
|
|
with self._sub_lock:
|
|
for code, ow in self._sub_owners.items():
|
|
if "_permanent" in ow:
|
|
keep.add(code)
|
|
for code, ow in self._us_sub_owners.items():
|
|
if "_permanent" in ow:
|
|
keep.add(code)
|
|
if get_env_bool("LS_WS_HOLD_RAM_KEEP_HOLDINGS", True):
|
|
fn = getattr(self, "_hold_ram_extra_codes_fn", None)
|
|
if callable(fn):
|
|
try:
|
|
extra = fn()
|
|
keep |= {str(c).strip() for c in (extra or set()) if str(c).strip()}
|
|
except Exception:
|
|
pass
|
|
return keep
|
|
|
|
def _force_close_socket(self, reason: str) -> None:
|
|
"""hold 이탈·갭 — 구독 UNREG 후 소켓 양보. 토큰 revoke/재발급 없음."""
|
|
ws = self._ws
|
|
if ws is None:
|
|
return
|
|
if self._force_closing:
|
|
return
|
|
self._force_closing = True
|
|
try:
|
|
try:
|
|
keep = self._hold_ram_keep_codes()
|
|
self._graceful_unreg_all(clear_ram=True, ram_keep_codes=keep)
|
|
except Exception as e:
|
|
logger.warning("LS WS 세션 양보 전 UNREG 실패: %s", e)
|
|
try:
|
|
logger.info("🔌 [LS WS] 세션 양보 close (%s)", reason)
|
|
ws.close()
|
|
except Exception as exc:
|
|
logger.debug("LS WS force close 예외: %s", exc)
|
|
finally:
|
|
self._force_closing = False
|
|
|
|
def _n_us_subscribed(self) -> int:
|
|
with self._sub_lock:
|
|
return len(self._us_subscribed)
|
|
|
|
def _should_hold_socket(self, now: Optional[datetime] = None) -> bool:
|
|
"""국내 hold 또는 (해외구독+해외 hold). ``LS_WS_SESSION_HOLD_ENABLED=false`` 면 항상 True."""
|
|
if not get_env_bool("LS_WS_SESSION_HOLD_ENABLED", True):
|
|
return True
|
|
from kis_trader.utils.ls_ws_session_windows import should_hold_ls_socket
|
|
|
|
return should_hold_ls_socket(
|
|
n_us_subscribed=self._n_us_subscribed(),
|
|
now=now,
|
|
)
|
|
|
|
def _session_guard_loop(self) -> None:
|
|
"""hold 창 이탈 시 소켓 능동 종료 — oauth 재발급 없음."""
|
|
from kis_trader.utils.ls_ws_session_windows import (
|
|
in_us_ws_hold_window,
|
|
session_guard_interval_sec,
|
|
)
|
|
|
|
while self._running:
|
|
try:
|
|
time.sleep(session_guard_interval_sec())
|
|
except Exception:
|
|
time.sleep(15.0)
|
|
if not self._running:
|
|
break
|
|
if not get_env_bool("LS_WS_SESSION_HOLD_ENABLED", True):
|
|
continue
|
|
if not self._opened.is_set():
|
|
continue
|
|
if self._should_hold_socket():
|
|
continue
|
|
why = "해외hold외" if not in_us_ws_hold_window() else "KR hold외"
|
|
self._force_close_socket("%s → 소켓 양보(토큰 유지)" % why)
|
|
|
|
def _in_kr_watchdog_wallclock(self, now: Optional[datetime] = None) -> bool:
|
|
"""국장 틱 silence 워치독용 정규장 벽시계 (기본 09:00~15:25, 월~금)."""
|
|
n = now or datetime.now()
|
|
if n.weekday() >= 5:
|
|
return False
|
|
hhmm = n.hour * 100 + n.minute
|
|
start = int(get_env_int("LS_WS_WATCHDOG_SESSION_START_HM", 900) or 900)
|
|
end = int(get_env_int("LS_WS_WATCHDOG_SESSION_END_HM", 1525) or 1525)
|
|
return start <= hhmm <= end
|
|
|
|
def _send_raw(self, payload: dict, *, close_on_fail: bool = True) -> None:
|
|
if self._ws is None or not self._opened.is_set():
|
|
return
|
|
timeout = max(0.5, float(get_env_float("LS_WS_SEND_LOCK_TIMEOUT_SEC", 3.0) or 3.0))
|
|
acquired = self._send_lock.acquire(timeout=timeout)
|
|
if not acquired:
|
|
logger.warning(
|
|
"LS WS send lock timeout %.1fs → UNREG 후 close (half-open 블로킹 방지)",
|
|
timeout,
|
|
)
|
|
if close_on_fail:
|
|
try:
|
|
self._force_close_socket("send lock timeout")
|
|
except Exception as e:
|
|
logger.debug("LS send-lock timeout close: %s", e)
|
|
return
|
|
send_failed = False
|
|
try:
|
|
try:
|
|
self._ws.send(json.dumps(payload, ensure_ascii=False))
|
|
except Exception as e:
|
|
logger.debug("LS WS send 실패: %s", e)
|
|
send_failed = True
|
|
finally:
|
|
self._send_lock.release()
|
|
if send_failed and close_on_fail:
|
|
try:
|
|
self._force_close_socket("send 실패")
|
|
except Exception as e2:
|
|
logger.debug("LS send 실패 후 close: %s", e2)
|
|
|
|
def _send_typed(
|
|
self, tr_type: str, tr_cd: str, tr_key: str, *, close_on_fail: bool = True,
|
|
) -> None:
|
|
self._send_raw(
|
|
{
|
|
"header": {"token": self._ensure_token(), "tr_type": str(tr_type)},
|
|
"body": {"tr_cd": tr_cd, "tr_key": tr_key},
|
|
},
|
|
close_on_fail=close_on_fail,
|
|
)
|
|
|
|
def _session_expects_trade_ticks(self, now: Optional[datetime] = None) -> bool:
|
|
"""틱 silence 강제재연결을 허용할 세션인가 (VI·장외 오탐 방지).
|
|
|
|
- ``LS_WS_WATCHDOG_SESSION_GATE=false`` 이면 항상 True(레거시).
|
|
- **벽시계 우선**: sticky JIF ``jstatus=21`` 만으로 장외 True 금지.
|
|
- 국장 구독: 정규창(기본 09:00~15:25) + 당일 마감 JIF/close 상태면 차단.
|
|
- 해외 구독: ``LS_WS_US_HOLD_*`` 창에서만.
|
|
"""
|
|
if not get_env_bool("LS_WS_WATCHDOG_SESSION_GATE", True):
|
|
return True
|
|
now = now or datetime.now()
|
|
ymd = now.strftime("%Y%m%d")
|
|
with self._sub_lock:
|
|
n_kr = len(self._subscribed)
|
|
n_us = len(self._us_subscribed)
|
|
kr_ok = False
|
|
if n_kr > 0 and self._in_kr_watchdog_wallclock(now):
|
|
close_set = _parse_jif_status_set(
|
|
"LS_AFR_JIF_CLOSE_STATUSES", _JIF_AFR_CLOSE_DEFAULT,
|
|
)
|
|
with self._jif_lock:
|
|
st = str(self._jstatus or "").strip()
|
|
closed = str(self._kr_afr_closed_ymd or "")
|
|
# 벽시계 안이면 기본 OK. sticky JIF 21 로 장외 연장 안 함.
|
|
# 당일 마감 JIF/close status 만 추가 차단.
|
|
if closed == ymd:
|
|
kr_ok = False
|
|
elif st and st in close_set:
|
|
kr_ok = False
|
|
else:
|
|
kr_ok = True
|
|
us_ok = False
|
|
if n_us > 0:
|
|
from kis_trader.utils.ls_ws_session_windows import in_us_ws_hold_window
|
|
|
|
us_ok = in_us_ws_hold_window(now)
|
|
return bool(kr_ok or us_ok)
|
|
|
|
def is_kr_afr_session_open(self) -> bool:
|
|
"""국장 LS 조건 AFR(t1860) 허용 세션인가.
|
|
|
|
해외 밤장과 무관. **장전 준비** + 마감 후 차단.
|
|
|
|
- 평일 ``LS_AFR_SESSION_START_HM``(기본 07:00)~END(15:30)
|
|
- **장전**: JIF 22~25 오기 전에도 START~``OPEN_DEADLINE``(기본 09:30) 전은 허용
|
|
→ AFR/스냅샷 미리 준비 (09:00 jstatus=21 기다리면 늦음)
|
|
- 개장 JIF(21~25) → sticky ON / opened_ymd
|
|
- 마감 JIF(30·31·41~44) → closed_ymd → 당일 재시도 금지
|
|
- DEADLINE 지나도 개장 JIF 없으면 휴장으로 보고 OFF (연타 방지)
|
|
"""
|
|
now = datetime.now()
|
|
if now.weekday() >= 5:
|
|
return False
|
|
hhmm = now.hour * 100 + now.minute
|
|
start = int(get_env_int("LS_AFR_SESSION_START_HM", 700) or 700)
|
|
end = int(get_env_int("LS_AFR_SESSION_END_HM", 1530) or 1530)
|
|
deadline = int(get_env_int("LS_AFR_OPEN_DEADLINE_HM", 930) or 930)
|
|
if not (start <= hhmm <= end):
|
|
return False
|
|
ymd = now.strftime("%Y%m%d")
|
|
with self._jif_lock:
|
|
live = bool(self._kr_afr_session_live)
|
|
opened = str(self._kr_afr_opened_ymd or "")
|
|
closed = str(self._kr_afr_closed_ymd or "")
|
|
# 당일 마감 후에는 상관없음(재시도 불필요)
|
|
if closed == ymd:
|
|
return False
|
|
# 당일 개장 신호 있으면 장중 유지 (jstatus 66/70 덮어쓰기 무시)
|
|
if live or opened == ymd:
|
|
return True
|
|
# 장전 준비 구간 — sticky 없어도 허용
|
|
if hhmm < deadline:
|
|
return True
|
|
# 장중 재기동: JIF sticky 유실(재시작 직후 jstatus 없음).
|
|
# 구로직은 DEADLINE 후 sticky 없으면 휴장으로 OFF → AFR 영구 0 + 주기 t1859
|
|
# 끄면 LS 유니버스 동결. 거래일·미마감이면 벽시계 장중 허용.
|
|
try:
|
|
from kis_trader.utils.kr_trading_day import is_kr_trading_day
|
|
|
|
if is_kr_trading_day(now.date()):
|
|
return True
|
|
except Exception:
|
|
pass
|
|
return False
|
|
|
|
def get_jif_snapshot(self) -> Dict[str, Any]:
|
|
"""진단용 — jangubun/jstatus/AFR sticky."""
|
|
with self._jif_lock:
|
|
return {
|
|
"jangubun": self._jangubun,
|
|
"jstatus": self._jstatus,
|
|
"jif_age_sec": (
|
|
(time.monotonic() - self._jif_mono) if self._jif_mono > 0 else None
|
|
),
|
|
"kr_afr_session_live": bool(self._kr_afr_session_live),
|
|
"kr_afr_session_src": self._kr_afr_session_src,
|
|
"kr_afr_opened_ymd": self._kr_afr_opened_ymd,
|
|
"kr_afr_closed_ymd": self._kr_afr_closed_ymd,
|
|
}
|
|
|
|
def _on_jif(self, body: Dict[str, Any]) -> None:
|
|
jangubun = str(body.get("jangubun") or "").strip()
|
|
jstatus = str(body.get("jstatus") or "").strip()
|
|
open_set = _parse_jif_status_set(
|
|
"LS_AFR_JIF_OPEN_STATUSES", _JIF_AFR_OPEN_DEFAULT,
|
|
)
|
|
close_set = _parse_jif_status_set(
|
|
"LS_AFR_JIF_CLOSE_STATUSES", _JIF_AFR_CLOSE_DEFAULT,
|
|
)
|
|
sticky_changed = False
|
|
ymd = datetime.now().strftime("%Y%m%d")
|
|
with self._jif_lock:
|
|
self._jangubun = jangubun
|
|
self._jstatus = jstatus
|
|
self._jif_mono = time.monotonic()
|
|
prev = bool(self._kr_afr_session_live)
|
|
if jstatus in open_set:
|
|
self._kr_afr_session_live = True
|
|
self._kr_afr_opened_ymd = ymd
|
|
# 다음날 개장 대비 — 당일 마감 플래그는 개장 시 클리어
|
|
if self._kr_afr_closed_ymd == ymd:
|
|
self._kr_afr_closed_ymd = ""
|
|
self._kr_afr_session_src = f"open:{jangubun}/{jstatus}"
|
|
elif jstatus in close_set:
|
|
self._kr_afr_session_live = False
|
|
self._kr_afr_closed_ymd = ymd
|
|
self._kr_afr_session_src = f"close:{jangubun}/{jstatus}"
|
|
sticky_changed = prev != bool(self._kr_afr_session_live)
|
|
live = bool(self._kr_afr_session_live)
|
|
logger.info(
|
|
"LS JIF 장운영 jangubun=%s jstatus=%s afr_live=%s%s",
|
|
jangubun, jstatus, live,
|
|
" (sticky변경)" if sticky_changed else "",
|
|
)
|
|
|
|
def _normalize_kr_code(self, raw: str) -> str:
|
|
code = str(raw or "").strip()
|
|
if code.startswith("U") and len(code) >= 7 and code[1:7].isdigit():
|
|
return code[1:7]
|
|
if code.isdigit() and len(code) == 6:
|
|
return code
|
|
# ex_shcode 등
|
|
digits = "".join(ch for ch in code if ch.isdigit())
|
|
if len(digits) >= 6:
|
|
return digits[-6:]
|
|
return code
|
|
|
|
def _on_hoga(self, tr_cd: str, body: Dict[str, Any]) -> None:
|
|
"""UH1(통합호가) 등 — RAM 갱신 + (옵션) interval 모드 DB 저장.
|
|
|
|
tick 모드(기본)는 체결 콜백에서 틱과 1:1 저장하므로 여기선 DB INSERT 안 함.
|
|
"""
|
|
if self._orderbook_cache is None:
|
|
return
|
|
code = self._normalize_kr_code(
|
|
str(body.get("shcode") or body.get("ex_shcode") or "")
|
|
)
|
|
if not (code.isdigit() and len(code) == 6):
|
|
return
|
|
src_map = {"UH1": "ls_uh1", "H1_": "ls_h1", "HA_": "ls_ha", "NH1": "ls_nh1"}
|
|
source = src_map.get(tr_cd, f"ls_{tr_cd.lower()}")
|
|
try:
|
|
snap = self._orderbook_cache.update_from_ls_hoga(code, body, source=source)
|
|
except Exception as e:
|
|
logger.debug("LS 호가 파싱 실패 %s: %s", code, e)
|
|
return
|
|
if snap is None:
|
|
return
|
|
if not self._orderbook_recorder:
|
|
return
|
|
if not get_env_bool("LS_WS_ORDERBOOK_SAVE", True):
|
|
return
|
|
mode = (
|
|
get_env_from_db("LS_WS_ORDERBOOK_SAVE_MODE", "tick") or "tick"
|
|
).strip().lower()
|
|
# tick / on_tick / tick_sync → 체결 경로에서만 저장 (여기는 RAM만)
|
|
if mode in ("tick", "on_tick", "tick_sync", "sync"):
|
|
return
|
|
gap_ms = max(0, int(get_env_int("LS_WS_ORDERBOOK_SAVE_MS", 1000) or 1000))
|
|
now_m = time.monotonic()
|
|
last = self._ob_last_save_mono.get(code, 0.0)
|
|
if gap_ms > 0 and (now_m - last) * 1000.0 < gap_ms:
|
|
return
|
|
self._ob_last_save_mono[code] = now_m
|
|
try:
|
|
self._orderbook_recorder(code, snap.to_storage_dict())
|
|
except Exception as e:
|
|
logger.debug("LS orderbook recorder: %s", e)
|
|
|
|
@staticmethod
|
|
def _vi_gubun_active(g: str) -> bool:
|
|
return str(g or "").strip() in ("1", "2", "3")
|
|
|
|
def _all_kr_in_vi(self) -> bool:
|
|
"""구독 중인 KR 이 1개 이상이고 전부 VI 활성인가."""
|
|
with self._sub_lock:
|
|
kr = [c for c in self._subscribed if c.isdigit() and len(c) == 6]
|
|
if not kr:
|
|
return False
|
|
return all(self.is_in_vi(c) for c in kr)
|
|
|
|
def _on_vi(self, tr_cd: str, body: Dict[str, Any]) -> None:
|
|
"""UVI/VI_ — RAM VI 상태 + ls_ws_vi 이벤트 적재."""
|
|
code = self._normalize_kr_code(
|
|
str(body.get("shcode") or body.get("ref_shcode") or body.get("ex_shcode") or "")
|
|
)
|
|
if not (code.isdigit() and len(code) == 6):
|
|
return
|
|
|
|
krx_g = str(body.get("krx_vi_gubun") or "").strip()
|
|
nxt_g = str(body.get("nxt_vi_gubun") or "").strip()
|
|
plain_g = str(body.get("vi_gubun") or "").strip()
|
|
if tr_cd == "UVI":
|
|
active = self._vi_gubun_active(krx_g) or self._vi_gubun_active(nxt_g)
|
|
# 저장용 대표 구분: KRX 우선, 없으면 NXT
|
|
vi_g = krx_g if krx_g != "" else (nxt_g if nxt_g != "" else "0")
|
|
if active and not self._vi_gubun_active(vi_g):
|
|
vi_g = nxt_g if self._vi_gubun_active(nxt_g) else krx_g
|
|
event_time = str(body.get("krx_time") or body.get("nxt_time") or "").strip()
|
|
try:
|
|
svi = float(body.get("krx_svi_recprice") or body.get("nxt_svi_recprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
svi = None
|
|
try:
|
|
dvi = float(body.get("krx_dvi_recprice") or body.get("nxt_dvi_recprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
dvi = None
|
|
try:
|
|
trg = float(body.get("krx_vi_trgprice") or body.get("nxt_vi_trgprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
trg = None
|
|
else:
|
|
active = self._vi_gubun_active(plain_g)
|
|
vi_g = plain_g or "0"
|
|
event_time = str(body.get("time") or "").strip()
|
|
try:
|
|
svi = float(body.get("svi_recprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
svi = None
|
|
try:
|
|
dvi = float(body.get("dvi_recprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
dvi = None
|
|
try:
|
|
trg = float(body.get("vi_trgprice") or 0) or None
|
|
except (TypeError, ValueError):
|
|
trg = None
|
|
krx_g = plain_g or None
|
|
nxt_g = None
|
|
|
|
now_m = time.monotonic()
|
|
with self._vi_lock:
|
|
if active:
|
|
self._vi_active_mono[code] = now_m
|
|
else:
|
|
self._vi_active_mono.pop(code, None)
|
|
|
|
logger.info(
|
|
"LS VI %s code=%s gubun=%s krx=%s nxt=%s active=%s",
|
|
tr_cd, code, vi_g, krx_g or "-", nxt_g or "-", active,
|
|
)
|
|
|
|
if active and get_env_bool("OPS_ALERT_VI_ENABLED", True):
|
|
try:
|
|
from kis_trader.utils.ops_alert import ops_alert
|
|
ops_alert(
|
|
f"vi_{code}",
|
|
f"VI 발동 {code}",
|
|
detail=(
|
|
f"tr={tr_cd} gubun={vi_g} krx={krx_g or '-'} "
|
|
f"nxt={nxt_g or '-'} trg={trg}"
|
|
),
|
|
level="warn",
|
|
session_only=False,
|
|
)
|
|
except Exception:
|
|
pass
|
|
|
|
if not self._vi_recorder or not get_env_bool("LS_WS_VI_SAVE", True):
|
|
return
|
|
payload = {
|
|
"ts": datetime.now(),
|
|
"event_time": event_time,
|
|
"vi_gubun": vi_g,
|
|
"krx_vi_gubun": krx_g,
|
|
"nxt_vi_gubun": nxt_g,
|
|
"svi_recprice": svi,
|
|
"dvi_recprice": dvi,
|
|
"vi_trgprice": trg,
|
|
"tr_cd": tr_cd,
|
|
"exchname": str(body.get("exchname") or ""),
|
|
"active": active,
|
|
}
|
|
try:
|
|
self._vi_recorder(code, payload)
|
|
except Exception as e:
|
|
logger.debug("LS VI recorder: %s", e)
|
|
|
|
def _reg_worker(self) -> None:
|
|
"""OPEN 콜백에서 sleep 하지 않고, 여기서 gap 두고 REG/UNREG."""
|
|
while self._running:
|
|
try:
|
|
item = self._reg_q.get(timeout=0.5)
|
|
except queue.Empty:
|
|
continue
|
|
if item is None:
|
|
break
|
|
kind = item[0]
|
|
gap_ms = max(20, int(get_env_int("LS_WS_REG_GAP_MS", 80) or 80))
|
|
if kind == "replay":
|
|
delay_ms = max(0, int(get_env_int("LS_WS_OPEN_REG_DELAY_MS", 500) or 0))
|
|
if delay_ms > 0:
|
|
time.sleep(delay_ms / 1000.0)
|
|
with self._sub_lock:
|
|
kr = list(self._subscribed)
|
|
us = list(self._us_subscribed)
|
|
n_ok = 0
|
|
# 관측 전용(동작 동일): 루프 break 원인·기대 send 수. 복구/재REG 로직 변경 없음.
|
|
abort_why = ""
|
|
also_hoga = bool(self.also_hoga)
|
|
uvi_on = get_env_bool("LS_WS_UVI_ENABLED", True)
|
|
jif_on = get_env_bool("LS_WS_JIF_ENABLED", True)
|
|
per_kr = 1 + (1 if also_hoga else 0) + (1 if uvi_on else 0)
|
|
expected = (1 if jif_on else 0) + len(kr) * per_kr + len(us)
|
|
# JIF 장운영 — 스펙 예: tr_key=0 (앱 하트비트 아님)
|
|
if jif_on:
|
|
jif_key = (get_env_from_db("LS_WS_JIF_TR_KEY", "0") or "0").strip() or "0"
|
|
self._send_typed("3", "JIF", jif_key)
|
|
n_ok += 1
|
|
time.sleep(gap_ms / 1000.0)
|
|
# replay 는 직접 전송 (큐 재투입 폭주 방지)
|
|
for code in kr:
|
|
if not self._running:
|
|
abort_why = "not_running"
|
|
break
|
|
if not self._opened.is_set():
|
|
abort_why = "opened_cleared"
|
|
break
|
|
tr_cd, tr_key, hoga_cd, hoga_key = self._kr_tr_pair(code)
|
|
self._send_typed("3", tr_cd, tr_key)
|
|
n_ok += 1
|
|
if also_hoga:
|
|
self._send_typed("3", hoga_cd, hoga_key)
|
|
n_ok += 1
|
|
if uvi_on:
|
|
vi_cd, vi_key = self._vi_tr_pair(code)
|
|
self._send_typed("3", vi_cd, vi_key)
|
|
n_ok += 1
|
|
time.sleep(gap_ms / 1000.0)
|
|
for sym in us:
|
|
if not self._running:
|
|
abort_why = abort_why or "not_running"
|
|
break
|
|
if not self._opened.is_set():
|
|
abort_why = abort_why or "opened_cleared"
|
|
break
|
|
self._send_typed("3", "GSC", overseas_tr_key("82", sym))
|
|
n_ok += 1
|
|
time.sleep(gap_ms / 1000.0)
|
|
ping_iv = int(get_env_int("LS_WS_PING_INTERVAL_SEC", 20) or 0)
|
|
reg_ok = (not abort_why) and self._opened.is_set() and n_ok >= expected
|
|
logger.info(
|
|
"LS WS OPEN — 구독 복구 KR=%d US=%d sends≈%d expected≈%d "
|
|
"abort_why=%s opened_now=%s also_hoga=%s reg_ok=%s "
|
|
"(delay=%dms gap=%dms ping=%s)",
|
|
len(kr),
|
|
len(us),
|
|
n_ok,
|
|
expected,
|
|
abort_why or "-",
|
|
self._opened.is_set(),
|
|
also_hoga,
|
|
reg_ok,
|
|
delay_ms,
|
|
gap_ms,
|
|
"off" if ping_iv <= 0 else f"{ping_iv}s",
|
|
)
|
|
# 미완료 REG 인데 recovering=False → 워치독/게이트가 정상으로 착각
|
|
if reg_ok:
|
|
self._mark_stable_open("REG replay complete")
|
|
self._set_recovering(False)
|
|
else:
|
|
self._set_recovering(True)
|
|
continue
|
|
# send
|
|
_, tr_type, tr_cd, tr_key = item
|
|
if self._opened.is_set():
|
|
self._send_typed(str(tr_type), str(tr_cd), str(tr_key))
|
|
time.sleep(gap_ms / 1000.0)
|
|
|
|
def _watchdog_loop(self) -> None:
|
|
"""정규장에서만: 구독 전체 틱 N초 없음 → 강제 close.
|
|
|
|
생존 1순위는 프로토콜 ping. 앱 JSON 하트비트는 보내지 않음.
|
|
장외·동시호가·마감·JIF 비정규 상태에서는 틱 silence로 끊지 않음.
|
|
"""
|
|
while self._running:
|
|
time.sleep(max(1, int(get_env_int("LS_WS_WATCHDOG_POLL_SEC", 5) or 5)))
|
|
if not get_env_bool("LS_WS_WATCHDOG_ENABLED", True):
|
|
continue
|
|
if not self._opened.is_set():
|
|
continue
|
|
with self._sub_lock:
|
|
n_sub = len(self._subscribed) + len(self._us_subscribed)
|
|
if n_sub <= 0:
|
|
continue
|
|
if not self._session_expects_trade_ticks():
|
|
continue
|
|
# 구독 KR 전부가 VI 중이면 틱 공백이 정상 → 강제재연결 금지
|
|
if get_env_bool("LS_WS_WATCHDOG_SKIP_WHEN_VI", True) and self._all_kr_in_vi():
|
|
continue
|
|
silence = max(15, int(get_env_int("LS_WS_WATCHDOG_SILENCE_SEC", 45) or 45))
|
|
grace = max(5, int(get_env_int("LS_WS_WATCHDOG_OPEN_GRACE_SEC", 30) or 30))
|
|
empty_after = max(1, int(get_env_int("LS_WS_WATCHDOG_EMPTY_BACKOFF_AFTER", 3) or 3))
|
|
empty_mult = max(2, int(get_env_int("LS_WS_WATCHDOG_EMPTY_BACKOFF_MULT", 4) or 4))
|
|
if self._watchdog_empty_streak >= empty_after:
|
|
silence = silence * empty_mult
|
|
now = time.monotonic()
|
|
if now - self._opened_mono < grace:
|
|
continue
|
|
if now - self._last_tick_mono < silence:
|
|
continue
|
|
if self.is_recovering():
|
|
continue
|
|
self._watchdog_empty_streak += 1
|
|
with self._jif_lock:
|
|
jst = self._jstatus or "-"
|
|
logger.warning(
|
|
"LS WS watchdog: %ds 틱 없음 (subs=%d streak=%d silence=%ds jstatus=%s) → 강제 재연결",
|
|
int(now - self._last_tick_mono),
|
|
n_sub,
|
|
self._watchdog_empty_streak,
|
|
silence,
|
|
jst,
|
|
)
|
|
try:
|
|
from kis_trader.utils.ops_alert import ops_alert
|
|
ops_alert(
|
|
"ws_tick_silence",
|
|
"LS 시세 WS 틱 공백 → 강제 재연결",
|
|
detail=(
|
|
f"silence={int(now - self._last_tick_mono)}s "
|
|
f"subs={n_sub} streak={self._watchdog_empty_streak} jstatus={jst}"
|
|
),
|
|
level="critical",
|
|
)
|
|
except Exception:
|
|
pass
|
|
self._set_recovering(True)
|
|
try:
|
|
keep = self._hold_ram_keep_codes()
|
|
self._graceful_unreg_all(clear_ram=True, ram_keep_codes=keep)
|
|
except Exception as e:
|
|
logger.debug("LS watchdog UNREG: %s", e)
|
|
try:
|
|
if self._ws is not None:
|
|
self._ws.close()
|
|
except Exception as e:
|
|
logger.debug("LS watchdog close: %s", e)
|
|
|
|
def _run_forever(self) -> None:
|
|
while self._running:
|
|
# ── hold 외 → 소켓 양보 후 다음 국내/해외 창까지 대기 (재발급 없음) ──
|
|
if not self._should_hold_socket():
|
|
if self._opened.is_set() or self._ws is not None:
|
|
self._force_close_socket("hold 외 — 루프 대기 전 양보")
|
|
from kis_trader.utils.ls_ws_session_windows import (
|
|
seconds_until_ls_socket_open,
|
|
)
|
|
|
|
wait_sec = float(
|
|
seconds_until_ls_socket_open(
|
|
n_us_subscribed=self._n_us_subscribed(),
|
|
)
|
|
)
|
|
logger.info(
|
|
"🌙 LS WS hold 외 — 재연결 중지, 다음 hold까지 %.0f분 대기 "
|
|
"(토큰 캐시 유지·재발급 없음)",
|
|
wait_sec / 60.0,
|
|
)
|
|
self._opened.clear()
|
|
self._set_recovering(True)
|
|
self._ws_reconnect_step = 0
|
|
self._watchdog_empty_streak = 0
|
|
for _ in range(int(wait_sec // 60)):
|
|
if not self._running:
|
|
return
|
|
time.sleep(60)
|
|
time.sleep(wait_sec % 60)
|
|
continue
|
|
|
|
try:
|
|
self._opened.clear()
|
|
self._set_recovering(True)
|
|
# 만료/임박만 갱신 — force 없음 (세션 전환·재연결이 공용 토큰 무효화 금지)
|
|
try:
|
|
self._ensure_token()
|
|
except Exception as e:
|
|
logger.warning("LS 토큰 ensure 실패(캐시/한도): %s", e)
|
|
self._ws_reconnect_step += 1
|
|
time.sleep(ws_reconnect_delay_for_attempt(self._ws_reconnect_step))
|
|
continue
|
|
self._ws = self._websocket.WebSocketApp(
|
|
self.ws_url,
|
|
on_open=self._on_open,
|
|
on_message=self._on_message,
|
|
on_error=self._on_error,
|
|
on_close=self._on_close,
|
|
)
|
|
# 헬퍼/문서와 동일: 프로토콜 ping 기본 ON(20).
|
|
# BAD_LENGTH 완화는 send lock + REG 워커(콜백 비서면)로 처리.
|
|
ping_iv = int(get_env_int("LS_WS_PING_INTERVAL_SEC", 20) or 0)
|
|
ping_to = int(get_env_int("LS_WS_PING_TIMEOUT_SEC", 10) or 10)
|
|
run_kw: Dict[str, Any] = {}
|
|
if ping_iv > 0:
|
|
run_kw["ping_interval"] = ping_iv
|
|
run_kw["ping_timeout"] = max(1, ping_to)
|
|
self._ws.run_forever(**run_kw)
|
|
except Exception as e:
|
|
logger.warning("LS WS run 예외: %s", e)
|
|
# 관측: run_forever 종료 후 clear (동작 동일)
|
|
logger.warning(
|
|
"LS WS run_forever 종료 → opened clear (ws_id=%s)",
|
|
id(self._ws) if self._ws is not None else None,
|
|
)
|
|
self._opened.clear()
|
|
self._set_recovering(True)
|
|
if not self._running:
|
|
break
|
|
# hold 이탈로 끊긴 경우 즉시 대기 분기로 (불필요 재발급·재연결 폭주 방지)
|
|
if not self._should_hold_socket():
|
|
continue
|
|
|
|
# 조기 Bye/CLOSE 서킷 — OPEN마다 step=0 리셋하던 버그 보정 + 토큰 1회 갱신
|
|
early = bool(self._last_close_early or self._pending_bye)
|
|
if early:
|
|
self._early_close_streak += 1
|
|
elif self._last_close_lived_sec >= self._stable_open_sec():
|
|
self._early_close_streak = 0
|
|
|
|
streak_need = max(2, int(get_env_int("LS_WS_EARLY_BYE_STREAK", 5) or 5))
|
|
if early and self._early_close_streak >= streak_need:
|
|
circuit_sleep = max(
|
|
15.0,
|
|
float(get_env_float("LS_WS_BYE_CIRCUIT_SLEEP_SEC", 60.0) or 60.0),
|
|
)
|
|
logger.warning(
|
|
"LS WS Bye 서킷 — early_streak=%d lived=%.1fs bye=%s "
|
|
"→ 토큰 무효화·재발급 후 %.0fs 대기",
|
|
self._early_close_streak,
|
|
self._last_close_lived_sec,
|
|
self._pending_bye,
|
|
circuit_sleep,
|
|
)
|
|
try:
|
|
self._refresh_token_after_bye()
|
|
except Exception as e:
|
|
logger.warning("LS Bye 서킷 토큰 갱신 실패: %s", e)
|
|
self._early_close_streak = 0
|
|
self._pending_bye = False
|
|
self._ws_reconnect_step = max(self._ws_reconnect_step, 1)
|
|
time.sleep(circuit_sleep)
|
|
continue
|
|
|
|
self._ws_reconnect_step += 1
|
|
time.sleep(ws_reconnect_delay_for_attempt(self._ws_reconnect_step))
|
|
|
|
def _on_open(self, _ws: Any) -> None:
|
|
# 콜백에서 sleep/REG 연타 금지 → 워커에 replay 위임
|
|
# ⚠️ flapping OPEN 에서 reconnect_step=0 리셋 금지 (Bye 루프 1초 연타 원인)
|
|
self._opened_mono = time.monotonic()
|
|
self._last_tick_mono = time.monotonic()
|
|
self._stable_open = False
|
|
self._last_close_early = False
|
|
self._opened.set()
|
|
self._set_recovering(True)
|
|
# 관측 전용: 소켓 신분(레이스 가설). clear/REG 동작은 그대로.
|
|
logger.info(
|
|
"LS WS OPEN cb ws_id=%s active_id=%s same=%s reconnect_step=%d early_streak=%d",
|
|
id(_ws) if _ws is not None else None,
|
|
id(self._ws) if self._ws is not None else None,
|
|
_ws is self._ws,
|
|
self._ws_reconnect_step,
|
|
self._early_close_streak,
|
|
)
|
|
self._enqueue_replay()
|
|
|
|
def _on_close(self, _ws: Any, status: Any, msg: Any) -> None:
|
|
# 관측 전용: 옛 소켓 close 여부. 동작은 기존과 동일하게 무조건 clear.
|
|
same = _ws is self._ws
|
|
msg_s = msg
|
|
try:
|
|
if isinstance(msg, (bytes, bytearray)):
|
|
msg_s = bytes(msg).decode("utf-8", errors="replace")
|
|
except Exception:
|
|
msg_s = repr(msg)
|
|
lived = 0.0
|
|
if self._opened_mono > 0:
|
|
lived = max(0.0, time.monotonic() - self._opened_mono)
|
|
stable = self._stable_open_sec()
|
|
msg_blob = str(msg_s or "")
|
|
is_bye = ("Bye" in msg_blob) or self._pending_bye
|
|
# 강제 세션 양보(hold 외)는 early Bye 서킷에 넣지 않음
|
|
early = (lived < stable) and (not self._force_closing)
|
|
self._last_close_lived_sec = lived
|
|
self._last_close_early = bool(early)
|
|
if is_bye:
|
|
self._pending_bye = True
|
|
logger.warning(
|
|
"LS WS CLOSE status=%s msg=%s ws_same=%s "
|
|
"close_id=%s active_id=%s opened_before=%s "
|
|
"lived=%.1fs early=%s bye=%s stable_was=%s",
|
|
status,
|
|
msg_s,
|
|
same,
|
|
id(_ws) if _ws is not None else None,
|
|
id(self._ws) if self._ws is not None else None,
|
|
self._opened.is_set(),
|
|
lived,
|
|
early,
|
|
is_bye,
|
|
self._stable_open,
|
|
)
|
|
self._stable_open = False
|
|
self._opened.clear()
|
|
self._set_recovering(True)
|
|
|
|
def _on_error(self, _ws: Any, err: Any) -> None:
|
|
err_s = ""
|
|
try:
|
|
if isinstance(err, (bytes, bytearray)):
|
|
err_s = bytes(err).decode("utf-8", errors="replace")
|
|
else:
|
|
err_s = str(err)
|
|
except Exception:
|
|
err_s = repr(err)
|
|
if "Bye" in err_s or "opcode=8" in err_s or "\\x03\\xe8" in err_s:
|
|
self._pending_bye = True
|
|
logger.warning(
|
|
"LS WS ERROR %s ws_same=%s err_id=%s active_id=%s bye=%s",
|
|
err,
|
|
_ws is self._ws,
|
|
id(_ws) if _ws is not None else None,
|
|
id(self._ws) if self._ws is not None else None,
|
|
self._pending_bye,
|
|
)
|
|
|
|
def _on_message(self, _ws: Any, message: Any) -> None:
|
|
try:
|
|
data = json.loads(message) if isinstance(message, str) else message
|
|
except Exception:
|
|
logger.warning(
|
|
"LS WS MSG non-json type=%s head=%r",
|
|
type(message).__name__,
|
|
(str(message)[:200] if message is not None else ""),
|
|
)
|
|
return
|
|
if not isinstance(data, dict):
|
|
logger.warning("LS WS MSG not-dict type=%s", type(data).__name__)
|
|
return
|
|
header = data.get("header") or {}
|
|
if not isinstance(header, dict):
|
|
header = {}
|
|
body_raw = data.get("body")
|
|
body: Dict[str, Any] = body_raw if isinstance(body_raw, dict) else {}
|
|
tr_cd = str(header.get("tr_cd") or "")
|
|
tr_key = str(header.get("tr_key") or "")
|
|
tr_type = str(header.get("tr_type") or "")
|
|
rsp_cd = str(
|
|
header.get("rsp_cd")
|
|
or header.get("rspCd")
|
|
or body.get("rsp_cd")
|
|
or ""
|
|
).strip()
|
|
rsp_msg = str(
|
|
header.get("rsp_msg")
|
|
or header.get("rspMsg")
|
|
or body.get("rsp_msg")
|
|
or ""
|
|
).strip()
|
|
# REG/UNREG 응답은 종종 body=null + header.rsp_cd/rsp_msg.
|
|
# 예전: body 없으면 즉시 return → 서버 거절 코드를 삼킴 (관측 공백).
|
|
if rsp_cd or rsp_msg or (
|
|
tr_type in ("1", "2", "3", "4") and not body
|
|
):
|
|
line = (
|
|
"LS WS RSP tr_type=%s tr_cd=%s tr_key=%s rsp_cd=%s rsp_msg=%s"
|
|
% (
|
|
tr_type or "-",
|
|
tr_cd or "-",
|
|
tr_key or "-",
|
|
rsp_cd or "-",
|
|
rsp_msg or "-",
|
|
)
|
|
)
|
|
if rsp_cd and rsp_cd not in ("00000", "0"):
|
|
logger.warning(line)
|
|
try:
|
|
from kis_trader.network.ls_token import is_ls_auth_error
|
|
|
|
if is_ls_auth_error(rsp_cd=rsp_cd, rsp_msg=rsp_msg):
|
|
self._pending_bye = True
|
|
logger.warning(
|
|
"LS WS auth RSP → Bye 서킷 후보 (rsp_cd=%s)",
|
|
rsp_cd,
|
|
)
|
|
except Exception:
|
|
pass
|
|
else:
|
|
logger.info(line)
|
|
if not body:
|
|
return
|
|
if not body:
|
|
return
|
|
if tr_cd == "JIF":
|
|
self._on_jif(body)
|
|
return
|
|
if tr_cd in ("UH1", "H1_", "HA_", "NH1"):
|
|
self._on_hoga(tr_cd, body)
|
|
return
|
|
if tr_cd in ("UVI", "VI_", "NVI", "DVI"):
|
|
self._on_vi(tr_cd, body)
|
|
return
|
|
code = str(body.get("shcode") or body.get("symbol") or "").strip()
|
|
if not code:
|
|
return
|
|
# U005930 → 005930
|
|
if code.startswith("U") and len(code) >= 7 and code[1:7].isdigit():
|
|
code = code[1:7]
|
|
price_raw = body.get("price")
|
|
try:
|
|
price = float(price_raw)
|
|
except (TypeError, ValueError):
|
|
return
|
|
if price <= 0:
|
|
return
|
|
now = time.time()
|
|
self._last_tick_mono = time.monotonic()
|
|
self._watchdog_empty_streak = 0
|
|
chetime = str(body.get("chetime") or body.get("kortm") or body.get("trdtm") or "")
|
|
cvol = body.get("cvolume") or body.get("trdq") or 0
|
|
totq = body.get("volume") or body.get("totq") or 0
|
|
try:
|
|
cvol_f = float(cvol or 0)
|
|
except (TypeError, ValueError):
|
|
cvol_f = 0.0
|
|
try:
|
|
tot_f = float(totq or 0)
|
|
except (TypeError, ValueError):
|
|
tot_f = 0.0
|
|
|
|
# KIS 캐시 호환 필드
|
|
row = {
|
|
"stck_prpr": str(int(price) if price >= 1000 else price),
|
|
"prdy_ctrt": str(body.get("drate") or body.get("rate") or ""),
|
|
"acml_vol": str(tot_f),
|
|
"cntg_vol": str(cvol_f),
|
|
"stck_oprc": str(body.get("open") or ""),
|
|
"stck_hgpr": str(body.get("high") or ""),
|
|
"stck_lwpr": str(body.get("low") or ""),
|
|
"chetime": chetime,
|
|
"_ts": now,
|
|
"_src": "ls",
|
|
"_tr_cd": tr_cd,
|
|
"_price_f": price,
|
|
}
|
|
skip_ram = False
|
|
try:
|
|
from kis_trader.engine.feed_fallback import (
|
|
is_feed_read_stale,
|
|
packet_lag_seconds,
|
|
)
|
|
skip_ram = is_feed_read_stale(packet_lag_seconds(chetime))
|
|
except Exception:
|
|
skip_ram = False
|
|
if not skip_ram:
|
|
with self._cache_lock:
|
|
self._cache[code] = row
|
|
self._notify_price_listeners(code, price, row)
|
|
|
|
# TICK_SAVE 는 ls_ws_ticks INSERT 만 게이트(_on_tick). 여기까지 막으면
|
|
# 호가 틱동기(LS_WS_ORDERBOOK_SAVE)도 같이 꺼진다. 늦은 체결시각은
|
|
# RAM skip 과 같이 콜백 생략(지금 호가를 옛 시각으로 찍으면 look-ahead).
|
|
if self._tick_recorder and not skip_ram:
|
|
try:
|
|
self._tick_recorder(
|
|
code,
|
|
{
|
|
"ts": datetime.now(),
|
|
"price": price,
|
|
"volume": cvol_f,
|
|
"tot_volume": tot_f,
|
|
"chetime": chetime,
|
|
"tr_cd": tr_cd,
|
|
},
|
|
)
|
|
except Exception as e:
|
|
logger.debug("LS tick recorder: %s", e)
|
|
|
|
if get_env_bool("LS_WS_CANDLE_SAVE", True) and not skip_ram:
|
|
fn = self._candle_save_filter
|
|
if fn is None or fn(code):
|
|
self._roll_candle(code, price, cvol_f, now)
|
|
|
|
def _roll_candle(self, code: str, price: float, cvol: float, now: float) -> None:
|
|
tf = max(1, get_env_int("LS_WS_CANDLE_TF_MIN", 1))
|
|
dt = datetime.fromtimestamp(now)
|
|
# 분 버킷
|
|
minute = (dt.minute // tf) * tf
|
|
bucket = dt.replace(minute=minute, second=0, microsecond=0)
|
|
key = bucket.strftime("%Y-%m-%d %H:%M:00")
|
|
flushed = None
|
|
with self._candle_lock:
|
|
cur = self._candles.get(code)
|
|
if cur and cur.get("datetime") != key:
|
|
flushed = dict(cur)
|
|
cur = None
|
|
if cur is None:
|
|
cur = {
|
|
"datetime": key,
|
|
"tf_min": tf,
|
|
"open": price,
|
|
"high": price,
|
|
"low": price,
|
|
"close": price,
|
|
"volume": max(0.0, cvol),
|
|
"tick_count": 1,
|
|
}
|
|
self._candles[code] = cur
|
|
else:
|
|
cur["high"] = max(float(cur["high"]), price)
|
|
cur["low"] = min(float(cur["low"]), price)
|
|
cur["close"] = price
|
|
cur["volume"] = float(cur.get("volume") or 0) + max(0.0, cvol)
|
|
cur["tick_count"] = int(cur.get("tick_count") or 0) + 1
|
|
if flushed:
|
|
# 확정봉 RAM (전략 get_candles)
|
|
try:
|
|
bar = self._forming_to_strategy_bar(flushed)
|
|
bar["is_confirmed"] = 1
|
|
bar["tf_min"] = int(flushed.get("tf_min") or tf)
|
|
self._push_confirmed(code, bar)
|
|
except Exception as e:
|
|
logger.debug("LS confirmed push: %s", e)
|
|
if self._candle_flusher:
|
|
try:
|
|
self._candle_flusher(code, flushed)
|
|
except Exception as e:
|
|
logger.debug("LS candle flush: %s", e)
|