""" kis_trader/ws/kis_ws.py — KIS WebSocket 실시간 체결가 캐시 (H0STCNT0) ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ 역할: check_sell_signals() 의 inquire_price() REST 폴링을 대체. WSManager가 구독한 종목(기본=후보·보유, MINIMAL=보유만) 체결 push → RAM 캐시. 영구구독 KR은 LS. 세션 41은 영구 전용 슬롯이 아님. WebSocket vs REST: REST : 요청마다 0.5s 딜레이(API 제한) + 응답 대기 → 보유 3종목 ≈ 3초 낭비 WebSocket: 연결 1회 + push 수신 → 체결 즉시 캐시 갱신, API 카운트 무관 KIS 공식 스펙: TR_ID : H0STCNT0 (국내주식 실시간체결가) 실전 URL : ws://ops.koreainvestment.com:21000 (KIS_WS_URL_REAL, env/DB 변경 가능) 모의 URL : ws://ops.koreainvestment.com:31000 (KIS_WS_URL_MOCK, env/DB 변경 가능) 세션 구독 한도: 최대 41종목 (H0STCNT0). 누가 들어가는지는 WSManager 라우팅. KIS 2026-02-24 경고 준수 (무한 재연결 차단 정책): - 재연결 횟수 제한: 1시간 내 MAX_RECONNECTS_PER_HOUR 회 초과 시 자동 대기 - 최대 총 재연결 횟수: MAX_RECONNECT_ATTEMPTS 회 초과 시 WebSocket 종료 → REST fallback - 지수 백오프 대신 단계 대기: env WS_RECONNECT_DELAY_SECS (기본 1,3,5,7,10초) - 정상 종료/세션 양보: 연결 종료 전 구독 해제(tr_type=2). 소켓만 닫지 않음. 모의투자(VTS): KIS 모의투자는 WebSocket을 지원하지 않는 경우가 많음. 연결 실패 시 is_active=False → check_sell_signals()가 REST로 자동 fallback. KIS_WS_MOCK_ENABLED=true (env/DB) 로 강제 활성화 가능. 설치 필요: pip install websocket-client """ from __future__ import annotations import json import logging import queue import random import threading import time from pathlib import Path from typing import Any, Dict, List, Optional, Set import requests logger = logging.getLogger("KISWebSocket") try: from .ws_reconnect_backoff import ws_reconnect_delay_for_attempt except ImportError: from kis_trader.ws.ws_reconnect_backoff import ws_reconnect_delay_for_attempt # ------------------------------------------------------------------ # 모듈 수준에서 get_env 함수를 참조 (kis_long_ver1 공용 함수 재사용) # 이 파일이 단독으로도 동작할 수 있도록 fallback import 포함 # ------------------------------------------------------------------ try: from kis_trader.utils.env import get_env_from_db, get_env_int, get_env_bool, get_env_float except ImportError: try: from kis_long_ver1 import get_env_from_db, get_env_int, get_env_bool 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_bool(key, default=False): # type: ignore[misc] return default def get_env_float(key, default): # type: ignore[misc] return default class KISWebSocketPriceCache: """ KIS H0STCNT0 실시간 체결가 WebSocket 수신기. 사용법: ws_cache = KISWebSocketPriceCache(app_key, app_secret, is_mock=False) ok = ws_cache.start() # 백그라운드 스레드 시작 ws_cache.subscribe("005930") # 종목 구독 data = ws_cache.get_price("005930") # inquire_price 호환 dict 반환 ws_cache.stop() get_price() 반환 값이 None 이면 → REST inquire_price() 로 fallback """ # H0STCNT0 데이터 필드 인덱스 ('^' 구분) # KIS 공식 columns (open-trading-api ccnl_krx / MCP 확인): # 10 ASKP1, 11 BIDP1, 12 CNTG_VOL(체결량), 13 ACML_VOL(누적거래량) # 과거 IDX_VOLUME=11 은 BIDP1(매수호가≈가격)을 읽어 ws_ticks.volume 이 오염됨. IDX_CODE = 0 # MKSC_SHRN_ISCD: 유가증권 단축 종목코드 IDX_TIME = 1 # STCK_CNTG_HOUR: 체결 시간 IDX_PRICE = 2 # STCK_PRPR: 주식 현재가 (체결가) IDX_SIGN = 3 # PRDY_VRSS_SIGN: 전일 대비 부호 IDX_CHANGE = 4 # PRDY_VRSS: 전일 대비 IDX_CHGPCT = 5 # PRDY_CTRT: 전일 대비율 # 당일 시고저 (매수 체크 시 REST inquire_price 대체용 → API 과부하 방지) IDX_OPEN = 7 # STCK_OPRC: 주식 시가 (당일) IDX_HIGH = 8 # STCK_HGPR: 주식 고가 (당일) IDX_LOW = 9 # STCK_LWPR: 주식 저가 (당일) IDX_ASKP1 = 10 # ASKP1: 매도호가1 IDX_BIDP1 = 11 # BIDP1: 매수호가1 IDX_CNTG_VOL = 12 # CNTG_VOL: 체결 거래량 (틱당, TickRecorder용) IDX_ACML_VOL = 13 # ACML_VOL: 누적 거래량 (봉 델타용) IDX_CTTR = 18 # CTTR: 체결강도 (ccnl_krx 공식 컬럼) IDX_BSOP_DATE = 33 # BSOP_DATE: 영업일자 YYYYMMDD IDX_VOLUME = 12 # 하위호환 alias → CNTG_VOL # ccnl_krx 공식 컬럼 수 (MKSC_SHRN_ISCD … VI_STND_PRC). 한 프레임 N건이면 이 폭으로 자른다. H0STCNT0_N_FIELDS = 46 # 재연결 정책 (KIS 2026-02-24 경고 준수) MAX_RECONNECT_ATTEMPTS = 10 # 총 재연결 최대 횟수 (STABLE_CONN_RESET_SEC 이상 안정 연결 후 끊기면 초기화) MAX_RECONNECTS_PER_HOUR = 6 # 1시간 내 재연결 허용 횟수 # 재연결 sleep 간격 — ws_reconnect_backoff (기본 1,3,5,7,10초) # 이 시간(초) 이상 안정적으로 연결이 유지됐다가 끊기면 카운터를 초기화. # 예: 5분 이상 정상 운영 후 네트워크 일시 장애 → '버스트 차단'이 아닌 '정상 재연결'로 간주. STABLE_CONN_RESET_SEC = 300.0 # 5분 # Approval key 유효시간 (23시간 캐시). KIS: access_token 갱신주기 6시간·유효 24시간. APPROVAL_KEY_CACHE_SEC = 82800 def __init__(self, app_key: str, app_secret: str, is_mock: bool = True, approval_slot: str = "main"): self.app_key = app_key self.app_secret = app_secret self.is_mock = is_mock slot = str(approval_slot or "main").strip().lower() self.approval_slot = slot if slot in ("main", "ob") else "main" # WebSocket URL (env/DB 로 재정의 가능 → 연결 실패 시 사용자가 수정) _default_real = "ws://ops.koreainvestment.com:21000" _default_mock = "ws://ops.koreainvestment.com:31000" self._ws_url = ( get_env_from_db("KIS_WS_URL_MOCK", _default_mock) if is_mock else get_env_from_db("KIS_WS_URL_REAL", _default_real) ) self._base_url = ( "https://openapivts.koreainvestment.com:29443" if is_mock else "https://openapi.koreainvestment.com:9443" ) # ── 가격 캐시 ────────────────────────────────────────────── # { code: {"data": dict, "ts": float} } self._cache: Dict[str, Dict] = {} self._cache_lock = threading.Lock() # 현재가 갱신 리스너 — BaseStrategy 틱매도 등 (캐시 lock 밖에서 호출) self._price_listeners: list = [] self._price_listener_lock = threading.Lock() # ── 구독 목록 ────────────────────────────────────────────── self._subscribed: Set[str] = set() self._sub_lock = threading.Lock() # WS 역할: "full"(시세+옵션호가) | "tick"(H0STCNT0만) | "orderbook"(H0STASP0만) # 매니저가 호가전용 2nd 세션을 띄우면 main=tick, ob=orderbook 으로 설정. self._ws_role: str = "full" # ── 영구 구독 목록 (홀딩 관심종목 등) ────────────────────── # unsubscribe() 호출에도 해제되지 않는 고정 구독 코드 self._permanent_codes: Set[str] = set() self._load_permanent_watchlist() # ── WebSocket 관련 ───────────────────────────────────────── self._ws = None # 현재 WebSocketApp 인스턴스 self._ws_thread: Optional[threading.Thread] = None self._running = False self._connected = False # ── approval_key (REST 토큰과 별개, WebSocket 전용 인증) ─── self._approval_key: Optional[str] = None self._approval_key_ts: float = 0.0 # invalid approval 연속 횟수(연결당 1회 집계) — N회면 6h 우회 응급 재발급 self._invalid_approval_streak: int = 0 self._invalid_approval_handled_this_conn: bool = False # ── 재연결 관리 ──────────────────────────────────────────── self._reconnect_count = 0 self._reconnect_times: list = [] # 최근 재연결 타임스탬프 목록 self._last_connect_time: float = 0.0 # 마지막 연결 성공 시각 (안정 연결 판단용) self._last_subscribe_error: str = "" # H0STCNT0 JSON rt_cd≠0 (즉시 끊김 원인 추적) self._subscribe_resend_lock = threading.Lock() # approval 세션 시간분할 — KR hold 이탈 시 능동 close (해외 WS 양보) self._session_guard_thread: Optional[threading.Thread] = None # ── CandleAggregator (스캘핑봇 연동 시 외부에서 주입) ───── # attach_candle_aggregator(agg) 로 연결, None이면 봉 집계 비활성 self._candle_agg: Optional["CandleAggregator"] = None # ── TickRecorder (C안: RAM 링버퍼 + ws_ticks 배치) ──────── self._tick_recorder: Optional[Any] = None # ── 호가 RAM (H0STASP0 · 2키 OB 세션) + TriggerSnapshotRecorder ── self._orderbook_cache: Optional[Any] = None self._trigger_snapshot_recorder: Optional[Any] = None # 메인 tick 세션이 호가 틱동기 시 참조할 2키 OB 세션 (kis_ws_ob) self._orderbook_peer: Optional[Any] = None try: from kis_trader.ws.orderbook_cache import OrderbookCache as _ObCache self._orderbook_cache = _ObCache() except Exception as _ob_ex: logger.warning("KIS OrderbookCache 초기화 실패: %s", _ob_ex) self._orderbook_cache = None # ── WS 연결 성공 시 갭보정 콜백 ─────────────────────────── # set_on_connected_callback(fn) 으로 등록. # 연결 성공(_on_open) 후 별도 스레드에서 호출됨. # 장 시간일 때만 실행 (새벽 재연결 시 API 빈 응답 방지) self._on_connected_callback: Optional[callable] = None # ── websocket-client 사용 가능 여부 ─────────────────────── try: import websocket as _ws_lib self._ws_lib = _ws_lib self._available = True except ImportError: self._ws_lib = None self._available = False logger.warning( "⚠️ websocket-client 미설치 → WebSocket 실시간 가격 비활성.\n" " pip install websocket-client 를 실행하세요." ) # ================================================================== # Public API # ================================================================== def attach_candle_aggregator(self, agg: "CandleAggregator") -> None: """ CandleAggregator 를 연결합니다. 연결 후부터 틱 수신 시 자동으로 agg.on_tick() 이 호출됩니다. kis_scalping_ver1.py 에서 ws_cache.attach_candle_aggregator(agg) 로 호출. """ self._candle_agg = agg logger.info("✅ CandleAggregator 연결 완료 (봉 집계 활성화)") def attach_tick_recorder(self, recorder: Any) -> None: """TickRecorder 연결 — H0STCNT0 틱 → RAM 링버퍼 + ws_ticks 배치.""" self._tick_recorder = recorder logger.info("✅ TickRecorder 연결 완료 (KIS H0STCNT0)") def attach_trigger_snapshot_recorder(self, recorder: Any) -> None: """TriggerSnapshotRecorder 연결 — H0STASP0 호가 → kis_ws_orderbook / filter_eval.""" self._trigger_snapshot_recorder = recorder logger.info("✅ TriggerSnapshotRecorder 연결 완료 (KIS H0STASP0)") def attach_orderbook_peer(self, peer: Any) -> None: """메인 tick 세션 → 2키 OB 세션 연결. 체결 1건당 호가 RAM 틱동기 DB용.""" self._orderbook_peer = peer logger.info("✅ KIS orderbook peer 연결 (틱동기 → 2키 OB RAM)") def set_on_connected_callback(self, fn) -> None: """ WS 연결 성공(_on_open) 시 호출할 콜백 등록. ScalpingBotV1._fill_all_gaps 를 넘겨서 매 연결 시 갭보정 자동 실행. 장 시간일 때만 콜백을 실행 (새벽 자동재연결 시 API 빈 응답 방지). """ self._on_connected_callback = fn def start(self, force_cleanup: bool = True) -> bool: """ WebSocket 수신 백그라운드 스레드를 시작합니다. 성공 시 True, 사용 불가(모의/패키지 미설치/키 없음) 시 False. False 를 반환해도 봇은 REST fallback으로 정상 동작합니다. """ if not self._available: return False # 모의투자는 기본 비활성 (KIS 모의 서버 WebSocket 미지원 가능성) # KIS_WS_MOCK_ENABLED=true 로 강제 활성화 가능 if self.is_mock and not get_env_bool("KIS_WS_MOCK_ENABLED", False): logger.info("ℹ️ 모의투자: WebSocket 기본 비활성 (KIS_WS_MOCK_ENABLED=true 로 활성 가능)") return False if self._running: return True # ── [CRITICAL] 비정상 종료 후 재시작 대비: 구독·가격 캐시만 정리 ────────────── # approval_key 는 .kis_approval_cache_*.json 파일 캐시 유지 (6h/24h KIS 정책). # 국내·해외 WS 가 동일 키 공유 — 재시작마다 REST 재발급하면 invalid approval 유발. if force_cleanup: logger.info("🧹 WebSocket 세션 초기화 (구독/가격 캐시 리셋, approval_key 파일 유지)") with self._sub_lock: self._subscribed.clear() with self._cache_lock: self._cache.clear() logger.info("✅ WebSocket 세션 초기화 완료 (구독/캐시 리셋)") # approval_key 발급 (연결 전 확인) if not self._get_approval_key(): logger.warning("⚠️ WebSocket approval_key 발급 실패 → REST fallback 모드") return False self._running = True self._ws_thread = threading.Thread( target=self._run_ws_loop, daemon=True, name="KIS-WS-H0STCNT0" ) self._ws_thread.start() self._session_guard_thread = threading.Thread( target=self._session_guard_loop, daemon=True, name="KIS-WS-KR-Guard" ) self._session_guard_thread.start() logger.info( "✅ KIS WebSocket 수신 스레드 시작 (H0STCNT0 | url=%s)", self._ws_url ) return True def _unsubscribe_all_before_close(self, reason: str = "", *, clear_ram: bool = False) -> int: """세션 close 전 서버 구독 전부 해제 (영구종목 포함). RAM 목록은 ``clear_ram=False`` 면 유지 — 장 양보 후 재접속 시 _on_open 이 다시 구독한다. stop() 만 RAM 을 비운다. """ if not self._ws: if clear_ram: with self._sub_lock: self._subscribed.clear() with self._cache_lock: self._cache.clear() return 0 with self._sub_lock: codes = sorted(self._subscribed) if clear_ram: self._subscribed.clear() if clear_ram: with self._cache_lock: self._cache.clear() if not codes: return 0 role = (self._ws_role or "full").strip().lower() gap = max(0.0, float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MIN_SEC", 0.08))) n = 0 for i, code in enumerate(codes): if i > 0 and gap > 0: time.sleep(gap) if role in ("full", "tick"): self._send_sub_msg(code, subscribe=False, tr_id="H0STCNT0") if role == "orderbook" or ( role == "full" and get_env_bool("WS_ORDERBOOK_SAVE_KIS", False) ): self._send_sub_msg(code, subscribe=False, tr_id="H0STASP0") n += 1 if gap > 0: time.sleep(gap) logger.info( "✅ 세션 종료 전 구독 해제 %d종목 role=%s (%s)", n, role, reason or "-", ) return n def stop(self, clear_subscriptions: bool = True) -> None: """ WebSocket 수신 중단 및 스레드 종료. 종료 전 서버 구독 해제(한투 정상 케이스). ``clear_subscriptions`` 는 호환용이며 False 여도 해제는 한다. """ try: self._unsubscribe_all_before_close(reason="stop", clear_ram=True) except Exception as e: logger.warning("종료 전 구독 해제 실패: %s", e) self._running = False self._connected = False if self._ws: try: self._ws.close() except Exception: pass if self._ws_thread and self._ws_thread.is_alive(): self._ws_thread.join(timeout=5) logger.info("🛑 KIS WebSocket 종료") # KIS WebSocket 세션 당 구독 가능 최대 종목 수 # 초과 시 서버에서 오류 반환 또는 계정 일시 차단 가능 (KIS 공지 준수) MAX_SUBSCRIPTIONS = 41 def _load_permanent_watchlist(self) -> None: """ long_term_watchlist.json 에 있는 홀딩 관심종목을 영구 구독 목록으로 로드. WS 연결 시 자동으로 subscribe(), unsubscribe() 호출 시에도 해제하지 않음. """ try: wl_path = Path(__file__).resolve().parent / "long_term_watchlist.json" if not wl_path.exists(): return items = json.loads(wl_path.read_text(encoding="utf-8")).get("items", []) codes = [i["code"] for i in items if i.get("code")] self._permanent_codes = set(codes) if codes: logger.info( "📌 홀딩 영구 구독 목록 로드: %s (%d종목)", ", ".join(codes), len(codes), ) except Exception as e: logger.warning("영구 구독 목록 로드 실패: %s", e) def subscribe(self, code: str) -> bool: """ 실시간 체결가 구독 등록. 이미 연결 중이면 즉시 구독 메시지 전송, 연결 전이면 연결 성공 시 일괄 등록. KIS 세션 한도(MAX_SUBSCRIPTIONS=41) 초과 시 등록 거부 후 False. 성공·이미구독 True / 빈코드·한도초과 False (매니저 spill 체인용). """ 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( "⚠️ WebSocket 구독 한도 초과(%d/%d) → %s 구독 거부 " "(KIS 세션 한도 준수: 불필요 종목 구독해제 후 재시도)", len(self._subscribed), self.MAX_SUBSCRIPTIONS, code, ) return False self._subscribed.add(code) if self._connected and self._ws: role = (self._ws_role or "full").strip().lower() if role in ("full", "tick"): self._send_sub_msg(code, subscribe=True, tr_id="H0STCNT0") if role == "orderbook" or ( role == "full" and get_env_bool("WS_ORDERBOOK_SAVE_KIS", False) ): self._send_sub_msg(code, subscribe=True, tr_id="H0STASP0") logger.info( "📡 WebSocket 구독 추가: %s (%d/%d) role=%s", code, len(self._subscribed), self.MAX_SUBSCRIPTIONS, self._ws_role, ) return True def unsubscribe(self, code: str) -> None: """ 실시간 체결가 구독 해제 및 캐시 삭제. 단, long_term_watchlist.json 의 영구 구독 종목은 해제하지 않음. """ code = (code or "").strip() if code in self._permanent_codes: logger.debug("📌 영구 구독 종목 해제 요청 무시: %s (홀딩 관심종목)", code) return with self._sub_lock: self._subscribed.discard(code) with self._cache_lock: self._cache.pop(code, None) if self._connected and self._ws: role = (self._ws_role or "full").strip().lower() if role in ("full", "tick"): self._send_sub_msg(code, subscribe=False, tr_id="H0STCNT0") if role == "orderbook" or ( role == "full" and get_env_bool("WS_ORDERBOOK_SAVE_KIS", False) ): self._send_sub_msg(code, subscribe=False, tr_id="H0STASP0") logger.info("📡 WebSocket 구독 해제: %s", code) 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]: """ WebSocket 캐시에서 실시간 가격 + 당일 시고저를 꺼냅니다. max_age_sec 초 이내 수신된 체결 틱만 유효. max_age_sec is None 이면 나이 무시(마지막 RAM — 실매 매수·매도). 체결 공백 ≠ 가격 삭제. 반환 형식: { "stck_prpr": "73900", # 현재가 (체결가, REST inquire_price와 동일) "prdy_vrss": "200", # 전일 대비 "prdy_ctrt": "0.27", # 전일 대비율 "stck_oprc": "73000", # 당일 시가 (REST와 동일) "stck_hgpr": "74500", # 당일 고가 (REST와 동일) "stck_lwpr": "72800", # 당일 저가 (REST와 동일) } 반환 None: WebSocket 미연결 / 데이터 없음 / max_age_sec 초 초과 → REST fallback 권장 ※ WS 캐시 정확도: H0STCNT0 체결 틱마다 KIS 서버가 시고저를 함께 전송 → 장중 REST inquire_price와 동일하며 지연이 없음 (REST 왕복 100~300ms 생략) """ if not self._available or not self._connected: return None with self._cache_lock: entry = self._cache.get(code) if not entry: return None age = time.time() - entry.get("ts", 0) if max_age_sec is not None and age > max_age_sec: # 데이터가 너무 오래됨 → REST로 재확인 return None data = entry.get("data") if isinstance(data, dict): out = dict(data) out["_age_ms"] = int(age * 1000) return out return data def get_orderbook(self, code: str, max_age_sec: float = 3.0) -> Optional[dict]: """메모리에 저장된 최신 KIS 호가(dict). OrderbookSnapshot → bid dict.""" if self._orderbook_cache is not None: try: return self._orderbook_cache.get_kis_bid_dict(code, max_age_sec=max_age_sec) except Exception: return None return None def get_orderbook_snapshot(self, code: str, max_age_sec: float = 3.0): """H0STASP0 RAM 호가 스냅샷 (실매 필터·폴백 체인용).""" if self._orderbook_cache is None: return None try: return self._orderbook_cache.get(code, max_age_sec=max_age_sec) except Exception: return None @property def is_active(self) -> bool: """WebSocket이 연결되어 실시간 데이터를 수신 중이면 True.""" return bool(self._available and self._connected and self._running) # ================================================================== # 내부 메서드 # ================================================================== def _approval_min_reissue_sec(self) -> float: """KIS access_token/approval REST 재발급 최소 간격(초). 공식: 갱신주기 6시간.""" return max(0.0, float(get_env_float("KIS_WS_APPROVAL_MIN_REISSUE_SEC", 21600.0))) def _approval_key_age_sec(self) -> float: if not self._approval_key_ts: return 999999.0 return max(0.0, time.time() - self._approval_key_ts) def _invalidate_approval_key(self, reason: str = "") -> None: """approval_key 캐시 클리어 — start(force_cleanup) 등 명시적 초기화에만 사용.""" had = bool(self._approval_key) self._approval_key = None self._approval_key_ts = 0.0 if had and reason: logger.info("🔑 WebSocket approval_key 무효화 (%s)", reason) def _subscribe_gap_sec(self) -> float: """종목별 H0STCNT0 구독 간격(초). 0에 가까우면 KIS 가 연결 직후 끊을 수 있음.""" lo = float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MIN_SEC", 0.08)) hi = float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MAX_SEC", 0.25)) if hi < lo: lo, hi = hi, lo if hi <= 0: return 0.0 return random.uniform(lo, hi) def _reconnect_session_wait_sec(self) -> float: """끊긴 직후 재접속 전 대기 — 서버 측 세션 해제 시간 확보.""" return max(0.0, float(get_env_float("KIS_WS_RECONNECT_SESSION_WAIT_SEC", 2.0))) def _instant_drop_cooldown_sec(self) -> float: """연속 즉시 끊김 후 재시도까지 대기(초). 장중·장외 분리.""" if not self._is_market_hours(): return self._seconds_until_market_open() return max( 30.0, float(get_env_float("KIS_WS_INSTANT_DROP_COOLDOWN_SEC", 90.0)), ) def _instant_drop_reason_text(self, streak: int, conn_sec: float) -> str: """즉시 끊김 로그용 — 장외 오진 방지, 실제 의심 원인 명시.""" parts: list[str] = [] if not self._is_market_hours(): parts.append("장외(WS 서비스 시간 외)") else: parts.append("장중") parts.append(f"연속 {streak}회 {conn_sec:.1f}초 내 종료") if self._last_subscribe_error: parts.append(f"구독오류={self._last_subscribe_error}") else: parts.append( "의심=①access_token 24h 만료·갱신(갱신주기 6h)과 겹침 " "②재연결마다 approval_key REST 재요청 " "③H0STCNT0 구독 연속 폭주(키움은 간격 있음)" ) return " | ".join(parts) def _get_approval_key( self, *, force_refresh: bool = False, bypass_min_reissue: bool = False, ) -> Optional[str]: """ WebSocket 전용 approval_key — kis_approval_manager 파일 캐시 (국내·해외 공유). KIS 정책: 24h 유효, 6h 이내 REST 재발급 금지, 재연결 시 동일 키 재사용. bypass_min_reissue: invalid approval 응급 시에만 True. """ try: from kis_approval_manager import KISApprovalManager except ImportError as exc: logger.error("kis_approval_manager import 실패: %s", exc) return None mgr = KISApprovalManager.instance(self.is_mock, slot=self.approval_slot) key = mgr.get_approval_key( self.app_key, self.app_secret, self._base_url, force_refresh=force_refresh, bypass_min_reissue=bypass_min_reissue, ) if key: self._approval_key = key self._approval_key_ts = mgr.issued_ts or time.time() return key def _emergency_reissue_approval(self, reason: str) -> Optional[str]: """ 서버가 캐시 approval 을 거부할 때 6h 가드 우회 REST 재발급. 기본 ON (KIS_WS_INVALID_APPROVAL_BYPASS_6H). 남용 금지 — invalid 연속만. """ if not get_env_bool("KIS_WS_INVALID_APPROVAL_BYPASS_6H", True): logger.warning( "🔑 invalid approval 응급 재발급 OFF (KIS_WS_INVALID_APPROVAL_BYPASS_6H=false)" ) return None try: from kis_approval_manager import KISApprovalManager except ImportError as exc: logger.error("kis_approval_manager import 실패: %s", exc) return None mgr = KISApprovalManager.instance(self.is_mock, slot=self.approval_slot) key = mgr.emergency_reissue( self.app_key, self.app_secret, self._base_url, reason=reason, ) if key: self._approval_key = key self._approval_key_ts = mgr.issued_ts or time.time() self._invalid_approval_streak = 0 logger.info( "✅ invalid approval 응급 재발급 완료 (앞8자: %s…)", key[:8], ) return key def _handle_invalid_approval(self, err_line: str) -> None: """ H0STCNT0 invalid approval — 파일 동기화 후, 연속 N회면 응급 REST 재발급. 연결당 1회만 처리 (종목별 구독 거부 폭주 방지). """ logger.warning("⚠️ H0STCNT0 구독 거부: %s", err_line) if self._invalid_approval_handled_this_conn: return self._invalid_approval_handled_this_conn = True self._invalid_approval_streak += 1 old_key = (self._approval_key or "")[:8] try: from kis_approval_manager import KISApprovalManager mgr = KISApprovalManager.instance(self.is_mock, slot=self.approval_slot) reloaded = mgr.reload_from_file() if reloaded and reloaded != (self._approval_key or ""): self._approval_key = reloaded self._approval_key_ts = mgr.issued_ts or self._approval_key_ts logger.info( "🔑 invalid approval → 파일에 다른 키 발견 (앞8자: %s…→%s…) 재접속 대기", old_key, reloaded[:8], ) return if reloaded: logger.info( "🔑 invalid approval → 파일 동일 키 (앞8자: %s…, streak=%d)", reloaded[:8], self._invalid_approval_streak, ) except Exception as exc: logger.debug("invalid approval 파일 동기화 실패: %s", exc) need_n = max(1, int(get_env_int("KIS_WS_INVALID_APPROVAL_REISSUE_AFTER", 2))) if self._invalid_approval_streak < need_n: logger.warning( "🔑 invalid approval streak %d/%d — 다음 끊김 후 재시도 " "(동일 키면 응급 재발급)", self._invalid_approval_streak, need_n, ) return new_key = self._emergency_reissue_approval( reason=f"invalid approval streak={self._invalid_approval_streak} {err_line}", ) if new_key and self._ws: try: self._unsubscribe_all_before_close( reason="invalid approval 재발급", clear_ram=False, ) except Exception as e: logger.warning("invalid approval 전 구독 해제 실패: %s", e) try: self._ws.close() except Exception: pass def _build_sub_payload(self, code: str, subscribe: bool, tr_id: str = "H0STCNT0") -> str: """구독(tr_type=1) / 해제(tr_type=2) JSON 메시지 생성.""" return json.dumps({ "header": { "approval_key": self._approval_key or "", "custtype": "P", "tr_type": "1" if subscribe else "2", "content-type": "utf-8", }, "body": { "input": { "tr_id": tr_id, "tr_key": code, } }, }) def _send_sub_msg(self, code: str, subscribe: bool = True, tr_id: str = "H0STCNT0") -> None: """WebSocket으로 구독/해제 메시지 전송. 실패 시 조용히 무시.""" if not self._ws: return try: self._ws.send(self._build_sub_payload(code, subscribe, tr_id)) except Exception as e: logger.debug("구독 메시지 전송 실패(%s): %s", code, e) def _parse_realtime_msg(self, raw: str) -> None: """ H0STCNT0 실시간 체결가 메시지 파싱 및 캐시 갱신. 정상 포맷: "0|H0STCNT0|001|005930^082317^73900^5^200^0.27^..." parts[0]: 암호화구분 (0=평문, 1=암호화) parts[1]: TR_ID parts[2]: 건수 parts[3]: 데이터 ('^' 구분) KIS PINGPONG: 메시지가 "PINGPONG" 문자열 → 동일하게 echoing. JSON 응답(구독 확인/에러): {"header":{...},"body":{...}} → 무시. """ if not raw: return # ── KIS Application-Level PINGPONG ───────────────────────── if raw.strip() == "PINGPONG": if self._ws: try: self._ws.send("PINGPONG") except Exception: pass return # ── 구독 응답(JSON) ───────────────────────────────────────── if raw.startswith("{"): try: j = json.loads(raw) header = j.get("header", {}) if header.get("tr_id") == "H0STCNT0": body = j.get("body", {}) rt = str(body.get("rt_cd", "") or "").strip() msg = str(body.get("msg1", "") or "").strip() tr_key = str((body.get("output") or {}).get("tr_key") or "").strip() if rt and rt != "0": err_line = f"rt_cd={rt} msg={msg}" + (f" code={tr_key}" if tr_key else "") self._last_subscribe_error = err_line if rt == "1" and "ALREADY IN SUBSCRIBE" in msg.upper(): logger.debug("H0STCNT0 구독 중복(무해): %s", err_line) elif "INVALID APPROVAL" in msg.upper(): self._handle_invalid_approval(err_line) else: logger.warning("⚠️ H0STCNT0 구독 거부: %s", err_line) else: # 구독 성공 → invalid streak 리셋 if self._invalid_approval_streak: self._invalid_approval_streak = 0 logger.debug( "H0STCNT0 구독 OK: %s", tr_key or msg or "SUCCESS", ) except Exception: pass return # ── 실시간 데이터 파싱 ────────────────────────────────────── parts = raw.split("|") if len(parts) < 4: return # parts[0]=암호화구분, parts[1]=TR_ID, parts[2]=건수, parts[3]=데이터 if parts[1] == "H0STASP0": # 실시간 호가 (H0STASP0) — RAM 스냅샷 + TriggerSnapshotRecorder try: raw_data = parts[3] fields = raw_data.split("^") if len(fields) < 43: return code_val = str(fields[0] or "").strip() if not code_val: return snap = None if self._orderbook_cache is not None: snap = self._orderbook_cache.update_from_kis_h0stasp0(code_val, fields) if snap is not None and self._trigger_snapshot_recorder is not None: try: self._trigger_snapshot_recorder.on_orderbook( snap, snap_time=getattr(snap, "snap_time", None) or None, ) except Exception as _rec_ex: logger.debug("H0STASP0 recorder 실패: %s", _rec_ex) # 레거시 DB 직접 INSERT (db 주입 시에만) if hasattr(self, "db") and self.db is not None and hasattr(self.db, "insert_kis_ws_orderbook"): try: legacy = { "BSOP_HOUR": fields[1] if len(fields) > 1 else "", "ASKP1": fields[3] if len(fields) > 3 else "", "BIDP1": fields[13] if len(fields) > 13 else "", "TOTAL_ASKP_RSQN": fields[43] if len(fields) > 43 else "", "TOTAL_BIDP_RSQN": fields[44] if len(fields) > 44 else "", } self.db.insert_kis_ws_orderbook(code=code_val, snap=legacy, market="KR") except Exception: pass except Exception as e: logger.debug("H0STASP0 호가 파싱 오류: %s", e) return if parts[1] != "H0STCNT0": return # 암호화된 데이터는 아직 미지원 (평문만 처리) if parts[0] == "1": logger.debug("H0STCNT0 암호화 데이터 수신 (처리 스킵) → REST fallback 권장") return # 한 메시지에 여러 건이 포함될 수 있음 (parts[2] = 건수). # 공식: 0|H0STCNT0|004|체결1^…^체결2^… → 4건 전부 RAM·봉·ws_ticks 에 넣는다. for fields in self._h0stcnt0_record_slices(parts[2], parts[3]): self._apply_h0stcnt0_one(fields) def _h0stcnt0_record_slices(self, count_raw: str, data: str) -> List[list]: """H0STCNT0 페이로드를 체결 1건씩 자른다. 폭이 안 맞으면 기존처럼 1건만.""" tokens = (data or "").split("^") min_need = max(self.IDX_PRICE, self.IDX_CHGPCT) if len(tokens) <= min_need: return [] try: n = int(str(count_raw or "1").strip() or "1") except (TypeError, ValueError): n = 1 n = max(1, n) width = int(self.H0STCNT0_N_FIELDS) if n > 1: if len(tokens) >= n * width: pass elif len(tokens) >= n and (len(tokens) % n) == 0: width = len(tokens) // n else: n = 1 if n <= 1: return [tokens] out: List[list] = [] for i in range(n): sl = tokens[i * width:(i + 1) * width] if len(sl) <= min_need: break out.append(sl) return out or [tokens] def _apply_h0stcnt0_one(self, fields: list) -> None: """H0STCNT0 체결 1건 → RAM 현재가 + 봉 + TickRecorder.""" if len(fields) <= max(self.IDX_PRICE, self.IDX_CHGPCT): return try: code = fields[self.IDX_CODE].strip() price = float(fields[self.IDX_PRICE]) if not code or price <= 0: return chg_raw = fields[self.IDX_CHANGE] if len(fields) > self.IDX_CHANGE else "0" pct_raw = fields[self.IDX_CHGPCT] if len(fields) > self.IDX_CHGPCT else "0.00" # 당일 시고저: REST inquire_price 대체용 (매수 체크 시 API 과부하 방지) open_raw = fields[self.IDX_OPEN] if len(fields) > self.IDX_OPEN else "0" high_raw = fields[self.IDX_HIGH] if len(fields) > self.IDX_HIGH else "0" low_raw = fields[self.IDX_LOW] if len(fields) > self.IDX_LOW else "0" # 체결시간·체결강도 (H0STCNT0 ccnl_krx — 주는 그대로 보존) cntg_hour_raw = ( str(fields[self.IDX_TIME]).strip() if len(fields) > self.IDX_TIME else "" ) bsop_date_raw = ( str(fields[self.IDX_BSOP_DATE]).strip() if len(fields) > self.IDX_BSOP_DATE else "" ) cntr_str_val: Optional[float] = None if len(fields) > self.IDX_CTTR: try: cntr_str_val = float( str(fields[self.IDX_CTTR]).strip().replace(",", "") or "0" ) except (ValueError, TypeError): cntr_str_val = None tick_time_raw: Optional[str] = None if bsop_date_raw and cntg_hour_raw: tick_time_raw = f"{bsop_date_raw}{cntg_hour_raw}" elif cntg_hour_raw: tick_time_raw = cntg_hour_raw # tick_time: 영업일+체결시각(초까지) 우선 — 없으면 HHMMSS만 if bsop_date_raw and cntg_hour_raw: tick_time = f"{bsop_date_raw}{cntg_hour_raw}" else: tick_time = cntg_hour_raw # inquire_price output 딕셔너리와 키 이름을 맞춤 # stck_oprc/hgpr/lwpr 도 함께 저장 → kis_short_ver2 매수 체크에서 REST 없이 활용 data_compat = { "stck_prpr": str(int(price)), # 현재가 (int 문자열) "prdy_vrss": chg_raw, # 전일 대비 "prdy_ctrt": pct_raw, # 전일 대비율 "stck_oprc": open_raw, # 당일 시가 "stck_hgpr": high_raw, # 당일 고가 "stck_lwpr": low_raw, # 당일 저가 } if cntr_str_val is not None: data_compat["cntr_str"] = str(cntr_str_val) if tick_time_raw: data_compat["kis_cntg_hour_raw"] = cntg_hour_raw data_compat["kis_bsop_date_raw"] = bsop_date_raw data_compat["tick_time"] = str(tick_time_raw) data_compat["chetime"] = str(cntg_hour_raw or "") # 읽기 RAM: 체결시각 vs 지금(2초). 적재(TickRecorder)는 버리지 않음. skip_ram = False try: from kis_trader.engine.feed_fallback import ( is_feed_read_stale, packet_lag_seconds, ) lag = packet_lag_seconds(tick_time_raw or tick_time) skip_ram = is_feed_read_stale(lag) except Exception: skip_ram = False if not skip_ram: with self._cache_lock: self._cache[code] = {"data": data_compat, "ts": time.time()} self._emit_price_listeners(code, float(price), data_compat) # 체결량·누적량 (MCP: CNTG_VOL=12, ACML_VOL=13) cntg_vol = 0 acml_vol = 0 if len(fields) > self.IDX_CNTG_VOL: try: cntg_vol = int(str(fields[self.IDX_CNTG_VOL]).strip() or "0") except ValueError: cntg_vol = 0 if len(fields) > self.IDX_ACML_VOL: try: acml_vol = int(str(fields[self.IDX_ACML_VOL]).strip() or "0") except ValueError: acml_vol = 0 # ── CandleAggregator: 키움증분 코드는 CNTG, KIS 전용은 ACML 델타 if self._candle_agg is not None and not skip_ram: use_inc = False try: use_inc = bool(self._candle_agg._volume_is_incremental(code)) except Exception: use_inc = False agg_vol = cntg_vol if use_inc else acml_vol self._candle_agg.on_tick(code, price, agg_vol, tick_time) # TickRecorder: 실매·백테 공통으로 틱당 체결량(CNTG) 저장 if self._tick_recorder is not None: try: self._tick_recorder.on_tick( code, price, cntg_vol, tick_time, source="kis", cntr_str=cntr_str_val, tick_time_raw=tick_time_raw, ) except Exception as ex: logger.debug("H0STCNT0→TickRecorder 실패 %s: %s", code, ex) # ── 호가 틱동기: 체결 1건당 2키 OB RAM 스냅 1장 (키움 0B↔0D 와 동일) # skip_ram 이면 밀린 체결시각에 지금 호가를 붙이지 않음(미래 호가 오염 방지). 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 ) peer = self._orderbook_peer if self._orderbook_peer is not None else self snap = None if peer is not None and hasattr(peer, "get_orderbook_snapshot"): snap = peer.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("H0STCNT0→호가틱동기 실패 %s: %s", code, ex) logger.debug("H0STCNT0 수신: %s → %s원", code, int(price)) except (ValueError, IndexError) as e: logger.debug("H0STCNT0 파싱 오류: %s", e) # ================================================================== # WebSocket 루프 (백그라운드 스레드) # ================================================================== def _is_market_hours(self) -> bool: """ 국내 WS 가 approval 세션을 점유해도 되는 시간 (KR hold). - env: KIS_WS_KR_HOLD_START_HM~END (기본 0700~2000, 장 전후·틱 적재 여유) - 모의: 서버가 09:00 이전 거부 → hold 안이어도 09:00 미만이면 False - 해외 hold 창과 겹치지 않게 설정 (갭에서 양쪽 close) """ import datetime as _dt from kis_trader.utils.kis_ws_session_windows import in_kr_ws_hold_window now = _dt.datetime.now() if not in_kr_ws_hold_window(now): return False if self.is_mock and now.time() < _dt.time(9, 0): return False return True def _seconds_until_market_open(self) -> float: """다음 국내 KR hold 시작까지 초.""" from kis_trader.utils.kis_ws_session_windows import seconds_until_kr_ws_open return float(seconds_until_kr_ws_open()) def _force_close_socket(self, reason: str) -> None: """run_forever 점유 중에도 세션 양보 — 구독 해제 후 close (approval 1:1).""" ws = self._ws if ws is None: return try: self._unsubscribe_all_before_close(reason=reason, clear_ram=False) except Exception as exc: logger.warning("세션 양보 전 구독 해제 실패: %s", exc) try: logger.info("🔌 [국내WS] 세션 양보 close (%s)", reason) ws.close() except Exception as exc: logger.debug("국내WS force close 예외: %s", exc) def _session_guard_loop(self) -> None: """KR hold 이탈·해외 hold 진입 시 국내 소켓 능동 종료.""" from kis_trader.utils.kis_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 self._connected: continue if not self._is_market_hours(): why = "해외hold" if in_us_ws_hold_window() else "KR hold외/장외" self._force_close_socket("%s → approval 세션 비움" % why) def _run_ws_loop(self) -> None: """ 재연결 정책을 포함한 WebSocket 메인 루프. run_forever() 가 반환(연결 끊김)되면 지수 백오프 후 재연결 시도. KIS 정책: 1시간 내 MAX_RECONNECTS_PER_HOUR 회 초과 시 강제 대기. KR hold 외(기본 20:00~07:00·주말·해외창): 재연결 없이 대기 + 소켓 close. """ # 즉시 끊김 감지 — env 로 조정 (기본: 3초 이내 × 3회) INSTANT_DROP_SEC = float(get_env_float("KIS_WS_INSTANT_DROP_SEC", 3.0)) INSTANT_DROP_MAX = int(get_env_int("KIS_WS_INSTANT_DROP_MAX", 3)) _instant_drop_streak = 0 _last_conn_duration = 0.0 while self._running: now = time.time() # ── KR hold 외 → 소켓 양보 후 다음 국내 창까지 대기 ──────── if not self._is_market_hours(): if self._connected: self._force_close_socket("KR hold 외 — 루프 대기 전 양보") wait_sec = self._seconds_until_market_open() logger.info( "🌙 국내 WS hold 외 — 재연결 중지, 다음 KR hold까지 %.0f분 대기 " "(approval 세션 양보·재발급 없음)", wait_sec / 60, ) # 재연결 카운터·백오프 초기화 (장 시작 시 깨끗하게 재접속) self._reconnect_count = 0 self._reconnect_times = [] _instant_drop_streak = 0 # 60초 단위로 쪼개서 슬립 (stop() 신호 빠르게 감지) for _ in range(int(wait_sec // 60)): if not self._running: return time.sleep(60) time.sleep(wait_sec % 60) continue # ── 연속 즉시 종료 → 쿨다운 (장외/장중 메시지 분리) ──────── if _instant_drop_streak >= INSTANT_DROP_MAX: wait_sec = self._instant_drop_cooldown_sec() reason = self._instant_drop_reason_text( _instant_drop_streak, _last_conn_duration, ) # invalid approval 루프면 쿨다운 전에 6h 우회 응급 재발급 if "INVALID APPROVAL" in (self._last_subscribe_error or "").upper(): need_n = max(1, int(get_env_int("KIS_WS_INVALID_APPROVAL_REISSUE_AFTER", 2))) if self._invalid_approval_streak >= need_n: self._emergency_reissue_approval( reason=f"instant_drop×{_instant_drop_streak} {reason}", ) if not self._is_market_hours(): logger.warning( "⚠️ KIS WebSocket 연속 즉시 끊김 %d회 — %s " "→ 다음 장 시작까지 %.0f분 대기", _instant_drop_streak, reason, wait_sec / 60, ) else: logger.warning( "⚠️ KIS WebSocket 연속 즉시 끊김 %d회 — %s " "→ %.0f초 쿨다운 후 재접속 " "(invalid면 응급 재발급 후 신규 approval_key 사용)", _instant_drop_streak, reason, wait_sec, ) self._reconnect_count = 0 self._reconnect_times = [] _instant_drop_streak = 0 self._last_subscribe_error = "" for _ in range(int(wait_sec // 60)): if not self._running: return time.sleep(60) time.sleep(wait_sec % 60) continue # ── 안정 연결 후 끊김이면 재연결 카운터 초기화 ─────────── # STABLE_CONN_RESET_SEC(5분) 이상 정상 운영 후 끊겼다면 # 이전 재연결 이력을 소거하여 'KIS 정상 케이스'로 처리. # (수개월 운영 중 산발적 네트워크 단절에 의해 10회 소진 → 영구 비활성화 방지) if ( self._last_connect_time > 0 and (now - self._last_connect_time) > self.STABLE_CONN_RESET_SEC and self._reconnect_count > 0 ): logger.info( "♻️ WebSocket 안정 운영(%.0f분) 후 끊김 → 재연결 카운터 초기화 (%d→0)", (now - self._last_connect_time) / 60, self._reconnect_count, ) self._reconnect_count = 0 self._reconnect_times = [] # ── 1시간 내 재연결 횟수 점검 ──────────────────────────── self._reconnect_times = [t for t in self._reconnect_times if now - t < 3600] if len(self._reconnect_times) >= self.MAX_RECONNECTS_PER_HOUR: # 가장 오래된 재연결로부터 1시간이 지날 때까지 대기 wait = max(10, 3600 - (now - self._reconnect_times[0]) + 10) logger.warning( "⛔ 1시간 내 WebSocket 재연결 %d회 초과 → %.0f초(%.0f분) 강제 대기 (KIS 차단 방지)", self.MAX_RECONNECTS_PER_HOUR, wait, wait / 60, ) time.sleep(wait) continue # ── 총 재연결 횟수 초과 → REST fallback 전환 ───────────── if self._reconnect_count >= self.MAX_RECONNECT_ATTEMPTS: logger.error( "❌ WebSocket 최대 재연결 %d회 초과 → WebSocket 종료, REST fallback 전환", self.MAX_RECONNECT_ATTEMPTS, ) self._running = False break # ── approval_key — 6h 이내 REST 재발급 금지, 재연결은 캐시 키 재사용 ── if self._reconnect_count > 0: sess_wait = self._reconnect_session_wait_sec() if sess_wait > 0: logger.debug( "⏳ WS 재접속 전 세션 대기 %.1fs (approval_key age=%.0f분)", sess_wait, self._approval_key_age_sec() / 60, ) time.sleep(sess_wait) approval_key = self._get_approval_key( force_refresh=get_env_bool("KIS_WS_RECONNECT_REFRESH_KEY", False), ) if not approval_key: _attempt = max(1, self._reconnect_count) _delay = ws_reconnect_delay_for_attempt(_attempt) logger.warning( "WebSocket approval_key 발급 불가 → %ds 대기 후 재시도", int(_delay), ) time.sleep(_delay) continue # ── 재연결 카운트 및 대기 ───────────────────────────────── self._reconnect_count += 1 if self._reconnect_count > 1: self._reconnect_times.append(time.time()) _attempt = self._reconnect_count - 1 _delay = ws_reconnect_delay_for_attempt(_attempt) logger.info( "🔄 WebSocket 재연결 시도 %d/%d (대기 %ds 완료)", self._reconnect_count, self.MAX_RECONNECT_ATTEMPTS, int(_delay), ) time.sleep(_delay) # ── WebSocketApp 생성 및 실행 ───────────────────────────── _conn_start = time.time() # 연결 시작 시각 (즉시 종료 감지용) try: ws_app = self._ws_lib.WebSocketApp( self._ws_url, on_open = self._on_open, on_message = self._on_message, on_error = self._on_error, on_close = self._on_close, ) self._ws = ws_app # ping_interval: WebSocket 프로토콜 레벨 PING (연결 유지) # KIS Application-Level PINGPONG은 _parse_realtime_msg 에서 처리 ws_app.run_forever(ping_interval=20, ping_timeout=10) except Exception as e: logger.error("WebSocket run_forever 예외: %s", e) if not self._running: break # ── 즉시 종료 여부 판정 ─────────────────────────────────────── _conn_duration = time.time() - _conn_start _last_conn_duration = _conn_duration if _conn_duration < INSTANT_DROP_SEC: _instant_drop_streak += 1 logger.debug( "⚡ WebSocket 즉시 종료 감지 (%.1f초, streak=%d/%d)", _conn_duration, _instant_drop_streak, INSTANT_DROP_MAX, ) else: # 어느 정도 연결 유지됐다면 즉시 종료 streak 초기화 _instant_drop_streak = 0 logger.info("KIS WebSocket 루프 종료 (is_active=False, REST fallback 전환)") # ------------------------------------------------------------------ # WebSocket 콜백 # ------------------------------------------------------------------ def _on_open(self, ws) -> None: """연결 성공: 등록된 모든 종목 구독 (간격 두고 1회만 — 이중 전송 금지).""" self._connected = True self._last_connect_time = time.time() self._last_subscribe_error = "" # 연결당 invalid approval 1회 처리 플래그 리셋 self._invalid_approval_handled_this_conn = False logger.info( "✅ KIS WebSocket 연결 성공 (H0STCNT0 | url=%s | approval_age=%.0f분)", self._ws_url, self._approval_key_age_sec() / 60, ) with self._sub_lock: for code in sorted(self._permanent_codes): if code not in self._subscribed: if len(self._subscribed) >= self.MAX_SUBSCRIPTIONS: logger.warning("⚠️ 구독 한도로 영구구독 추가 불가: %s", code) break self._subscribed.add(code) codes = sorted(self._subscribed) def _subscribe_all() -> None: role = (self._ws_role or "full").strip().lower() for i, code in enumerate(codes): if i > 0: gap = self._subscribe_gap_sec() if gap > 0: time.sleep(gap) if role in ("full", "tick"): self._send_sub_msg(code, subscribe=True, tr_id="H0STCNT0") if role == "orderbook" or ( role == "full" and get_env_bool("WS_ORDERBOOK_SAVE_KIS", False) ): self._send_sub_msg(code, subscribe=True, tr_id="H0STASP0") if codes: logger.info( "📡 WebSocket 구독 일괄 등록: %s (%d종목, gap=%.2f~%.2fs, role=%s)", ", ".join(codes), len(codes), float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MIN_SEC", 0.08)), float(get_env_float("KIS_WS_SUBSCRIBE_GAP_MAX_SEC", 0.25)), role, ) if codes: threading.Thread( target=_subscribe_all, daemon=True, name="KIS-WS-Sub", ).start() # ── 연결 성공 시 갭보정 콜백 (장 시간일 때만) ─────────────── # 봇 시작 후 처음 WS가 안정 연결되는 시점(9:00 이후)에 # REST 분봉 데이터로 DB를 채워 '봉부족' 문제 해소. # 새벽 재연결 시에는 is_market_hours() False → 콜백 스킵. if self._on_connected_callback and self._is_market_hours(): try: cb_thread = threading.Thread( target=self._on_connected_callback, name="WS-GapFill", daemon=True, ) cb_thread.start() except Exception as e: logger.debug("갭보정 콜백 실행 실패: %s", e) def _on_message(self, ws, message: str) -> None: self._parse_realtime_msg(message) def _on_error(self, ws, error) -> None: self._connected = False logger.warning("⚠️ KIS WebSocket 오류: %s", error) def _on_close(self, ws, close_status_code, close_msg) -> None: self._connected = False logger.info( "🔌 KIS WebSocket 연결 종료 (code=%s msg=%s)", close_status_code, close_msg or "", ) # ====================================================================== # ====================================================================== # CandleAggregator — WebSocket 틱 → OHLCV 봉 실시간 집계기 # ====================================================================== class CandleAggregator: """ KISWebSocketPriceCache 에서 수신한 틱을 N분봉으로 집계합니다. Two-Track 아키텍처 (Gemini/퀀트 펌 방식) ───────────────────────────────────────── [트랙 1 — 매매 두뇌 (논블로킹)] WebSocket 틱 → on_tick() → RAM에서 OHLCV 즉시 갱신 → 매수/매도 판단은 get_latest_confirmed() 등 메모리 접근만 사용 → DB 대기 시간 0ms, 타점 놓침 없음 [트랙 2 — 기록원 스레드 (백그라운드)] 봉 확정 시 dict를 Queue에 put_nowait() (논블로킹, 0.000001초) → 백그라운드 스레드(_db_writer)가 BATCH_SIZE개 or FLUSH_INTERVAL초마다 DB에 executemany() 한 방에 묶어서 INSERT → 백테스트용 봉 데이터 완전 보존, 매매 루프 블로킹 없음 진행 중 봉(is_confirmed=0) 처리 ───────────────────────────────── - RAM의 _current 버퍼에만 존재 → DB에 절대 쓰지 않음 - 매수 루프는 get_current_candle() 로 즉시 메모리 접근 - 봉 확정(분 바뀜) 순간에만 Queue → DB 기록 RSI 계산 ──────── - 확정 봉 close 리스트(_closes, 최대 MAX_CLOSE_BUFFER개) RAM에 유지 - RSI(2/3/5) 계산은 순수 Python 연산, DB 조회 없음 - 봉 수 부족 시 RSI=None → 매수 루프에서 신호 무시 스레드 안전성 ───────────── - _lock : on_tick / fill_gap 간 경합 방지 (RAM 버퍼 보호) - Queue : thread-safe, put_nowait 는 lock 불필요 - _db_writer : 독립 daemon 스레드 (봇 종료 시 자동 소멸) """ MAX_CLOSE_BUFFER = 200 # RSI 계산용 close 보관 기본값 (WS_CANDLE_RAM_BUFFER 로 덮어씀) BATCH_SIZE = 50 # 이 개수 이상 쌓이면 즉시 배치 플러시 FLUSH_INTERVAL = 2.0 # 초 — BATCH_SIZE 미달이라도 이 주기로 플러시 def __init__(self, db=None, timeframes: list = None): """ Args: db : TradeDB 인스턴스. None이면 DB 쓰기 비활성(순수 RAM 모드). timeframes: 집계할 봉 단위 리스트 (기본 [1, 3]) """ self.db = db self.timeframes: list = timeframes if timeframes else [1, 3] # MOMENTUM 전일시가(E) 등 — RAM에 최소 2영업일 1분봉 보관 (DB 병합 없이 키움 REST 갭보정) self._ram_buffer_max = max( self.MAX_CLOSE_BUFFER, get_env_int("WS_CANDLE_RAM_BUFFER", 500), ) self._lock = threading.Lock() # ── 트랙 1: RAM 버퍼 ───────────────────────────────────── # 진행 중인 봉: { (code, tf): {candle_time, open, high, low, close, volume} } self._current: Dict[tuple, dict] = {} # RSI 계산용 확정 봉 close 가격 히스토리: { (code, tf): [c1, c2, ...] } self._closes: Dict[tuple, list] = {} # 최근 확정 봉 보관 (get_latest_confirmed / get_candles 용): { (code, tf): [dict, ...] } self._confirmed: Dict[tuple, list] = {} # ── 트랙 2: 비동기 DB 배치 큐 ──────────────────────────── # 봉 확정 시 dict를 여기 넣으면 _db_writer 가 배치로 INSERT self._write_queue: queue.Queue = queue.Queue(maxsize=10000) self._db_writer_thread: Optional[threading.Thread] = None # 보유 중 트레일 고점: code → float (SHORT 봇이 set_peak_provider 로 주입) self._peak_provider: Optional[Any] = None # 키움 WS 틱 체결량(FID 15) → 분봉 volume 합산. 미등록 종목은 KIS 누적거래량 델타. self._incremental_vol_codes: Set[str] = set() self._tick_lookup: Optional[Any] = None self._ls_bars_fn: Optional[Any] = None # 핫패스(락 안)에서 get_env 금지 — env 스냅샷 cold(수 초)가 락을 잡으면 # 전 전략 get_candles 가 멈춘다. db_writer 가 주기 갱신. self._freeze_on_confirm_cached: bool = bool( get_env_bool("WS_CANDLE_FREEZE_ON_CONFIRM", True) ) self._flags_refresh_ts: float = 0.0 if self.db is not None: self._start_db_writer() logger.info( "✅ CandleAggregator 초기화 완료 (timeframes=%s, db=%s)", self.timeframes, "활성" if self.db else "비활성(RAM 전용)", ) def set_peak_provider(self, provider) -> None: """보유 종목의 트레일링 고점 조회 — 봉 확정 시 ws_candles.holding_peak 저장용.""" self._peak_provider = provider def set_incremental_volume_codes(self, codes: Optional[Set[str]]) -> None: """키움 WS 후보 종목 — 틱별 체결량을 분봉 volume 에 누적(+=).""" with self._lock: if codes is None: self._incremental_vol_codes = set() else: self._incremental_vol_codes = {str(c).strip() for c in codes if c} def set_tick_lookup(self, fn: Optional[Any]) -> None: """실매 쓰레기 봉 판정용 TickRecorder 조회. fn(code) -> list[dict].""" self._tick_lookup = fn def set_ls_bars_fn(self, fn: Optional[Any]) -> None: """3차 LS 확정봉 RAM. fn(code, tf) -> list[dict] source=ls.""" self._ls_bars_fn = fn def _volume_is_incremental(self, code: str) -> bool: return str(code).strip() in self._incremental_vol_codes def _holding_peak_for(self, code: str) -> Optional[float]: if not self._peak_provider: return None try: v = self._peak_provider(code) return float(v) if v and float(v) > 0 else None except Exception: return None # ------------------------------------------------------------------ # 트랙 2: 백그라운드 DB 배치 기록원 # ------------------------------------------------------------------ def _start_db_writer(self) -> None: """백그라운드 DB 배치 기록 스레드를 시작합니다.""" self._db_writer_thread = threading.Thread( target=self._db_writer_loop, name="CandleDBWriter", daemon=True, # 봇 종료 시 자동 소멸 ) self._db_writer_thread.start() logger.info("✅ CandleAggregator DB 기록원 스레드 시작 (배치=%d, 주기=%.1fs)", self.BATCH_SIZE, self.FLUSH_INTERVAL) def _db_writer_loop(self) -> None: """ 배치 기록 루프. - Queue에서 BATCH_SIZE개 모이면 즉시 배치 INSERT - BATCH_SIZE 미달이라도 FLUSH_INTERVAL초마다 플러시 - 겸사겸사 봉주기 경과 후 미확정 상태로 방치된 저유동 종목 봉도 같은 주기로 강제확정 점검 (flush_stale_current_candles) """ batch: list = [] last_flush = time.time() last_stale_check = 0.0 stale_check_interval = float( get_env_int("WS_CANDLE_STALE_CHECK_INTERVAL_SEC", 2) ) while True: try: # FLUSH_INTERVAL 내에서 최대한 많이 모아 배치 구성 timeout = max(0.1, self.FLUSH_INTERVAL - (time.time() - last_flush)) item = self._write_queue.get(timeout=timeout) if item is None: # None = 종료 신호 break batch.append(item) self._write_queue.task_done() except queue.Empty: pass now = time.time() should_flush = len(batch) >= self.BATCH_SIZE or \ (batch and (now - last_flush) >= self.FLUSH_INTERVAL) if should_flush and batch: self._flush_batch(batch) batch = [] last_flush = now if now - last_stale_check >= stale_check_interval: last_stale_check = now # 핫리로드 — 재시작 없이 WS_CANDLE_STALE_CHECK_INTERVAL_SEC 반영 try: stale_check_interval = max( 0.2, float(get_env_int("WS_CANDLE_STALE_CHECK_INTERVAL_SEC", 2) or 2), ) except Exception: pass # 락 밖: env 핫플래그 갱신 (on_tick 경로에서 get_env 금지) self._refresh_runtime_flags() try: self.flush_stale_current_candles() except Exception as e: logger.debug("flush_stale_current_candles 실패(무시): %s", e) # 루프 종료 시 남은 배치 처리 if batch: self._flush_batch(batch) logger.info("CandleDBWriter 스레드 종료") def _refresh_runtime_flags(self) -> None: """락 밖에서만 호출. freeze 핫패스 플래그 캐시 갱신.""" try: self._freeze_on_confirm_cached = bool( get_env_bool("WS_CANDLE_FREEZE_ON_CONFIRM", True) ) self._flags_refresh_ts = time.time() except Exception: pass def _ws_candle_freeze_on_confirm(self) -> bool: """확정봉 OHLCV 동결 — docs/정합성.md. 기본 true. ⚠ on_tick/_confirm 은 candle_agg._lock 보유 중 호출됨. 여기서 get_env_* 하면 ENV 캐시 만료 시 DB 스냅샷(수 초)이 락을 잡아 전 전략 매수루프가 멈춘다 → 캐시만 반환. """ return bool(getattr(self, "_freeze_on_confirm_cached", True)) @staticmethod def _bar_key(code: str, tf: int, source: str = "kis", channel: str = "") -> tuple: from kis_trader.ws.candle_series import normalize_source_channel src, ch = normalize_source_channel(source, channel) return (code, int(tf), src, ch) def _load_confirmed_ohlcv_from_db( self, code: str, tf: int, candle_times: list, source: str = "kis", channel: str = "", ) -> Dict[str, Dict]: """ freeze 재시작 정합: RAM 이 비어도 DB 에 이미 확정된 봉은 REST 로 덮지 않고 DB 값을 RAM 에 시드한다. (실매 RAM ≠ DB 재발 방지) """ out: Dict[str, Dict] = {} if not self.db or not candle_times: return out times = sorted({str(t)[:12] for t in candle_times if str(t)[:12]}) if not times: return out from kis_trader.ws.candle_series import normalize_source_channel src, ch = normalize_source_channel(source, channel) chunk_n = max(50, int(get_env_int("WS_CANDLE_FREEZE_DB_LOOKUP_CHUNK", 200))) try: for i in range(0, len(times), chunk_n): chunk = times[i : i + chunk_n] ph = ",".join(["%s"] * len(chunk)) rows = self.db.conn.execute( f""" SELECT candle_time, `open`, high, low, close, volume, rsi_2, rsi_3, rsi_5, source, channel, holding_peak FROM ws_candles WHERE code=%s AND timeframe=%s AND source=%s AND channel=%s AND is_confirmed=1 AND candle_time IN ({ph}) """, (code, int(tf), src, ch, *chunk), ).fetchall() for r in rows or []: ct = str(r.get("candle_time") or "")[:12] if not ct: continue out[ct] = dict(r) except Exception as e: logger.debug("freeze DB lookup 실패(%s %dM): %s", code, tf, e) return out @staticmethod def _ws_candles_upsert_sql(*, freeze: bool) -> str: """ freeze ON: 이미 확정된 행의 OHLCV/RSI/volume/source 유지. 미확정(is_confirmed=0)→확정·갱신은 허용. holding_peak 만 항상 GREATEST. freeze OFF: 기존처럼 덮어쓰기. """ if freeze: dup = """ ON DUPLICATE KEY UPDATE market=IF(VALUES(market) IS NULL OR VALUES(market)='', market, VALUES(market)), `open`=IF(is_confirmed=1, `open`, VALUES(`open`)), high=IF(is_confirmed=1, high, GREATEST(high, VALUES(high))), low=IF(is_confirmed=1, low, VALUES(low)), close=IF(is_confirmed=1, close, VALUES(close)), volume=IF(is_confirmed=1, volume, VALUES(volume)), rsi_2=IF(is_confirmed=1, rsi_2, VALUES(rsi_2)), rsi_3=IF(is_confirmed=1, rsi_3, VALUES(rsi_3)), rsi_5=IF(is_confirmed=1, rsi_5, VALUES(rsi_5)), is_confirmed=IF(is_confirmed=1, 1, VALUES(is_confirmed)), source=IF(is_confirmed=1, source, VALUES(source)), channel=IF(is_confirmed=1, channel, VALUES(channel)), holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)), updated_at=IF(is_confirmed=1, updated_at, VALUES(updated_at)) """ else: dup = """ ON DUPLICATE KEY UPDATE market=IF(VALUES(market) IS NULL OR VALUES(market)='', market, VALUES(market)), `open`=VALUES(`open`), high=GREATEST(high, VALUES(high)), low=VALUES(low), close=VALUES(close), volume=VALUES(volume), rsi_2=VALUES(rsi_2), rsi_3=VALUES(rsi_3), rsi_5=VALUES(rsi_5), is_confirmed=VALUES(is_confirmed), source=VALUES(source), channel=VALUES(channel), holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)), updated_at=VALUES(updated_at) """ return f""" INSERT INTO ws_candles (code, market, timeframe, candle_time, `open`, high, low, close, volume, rsi_2, rsi_3, rsi_5, is_confirmed, source, channel, holding_peak, updated_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) {dup} """ def _flush_batch(self, batch: list) -> None: """ 배치 리스트를 DB에 한 번의 executemany 로 INSERT. 각 item: {code, market, tf, candle_time, open, high, low, close, volume, is_confirmed, source, channel, rsi_2, rsi_3, rsi_5} """ if not self.db or not batch: return try: import datetime as _dt now_str = _dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S") rows = [] from kis_trader.ws.candle_series import normalize_source_channel for item in batch: src, ch = normalize_source_channel( item.get("source"), item.get("channel") or "", ) rows.append(( item["code"], str(item.get("market") or "KR").upper()[:8], item["tf"], item["candle_time"], item["open"], item["high"], item["low"], item["close"], item.get("volume", 0), item.get("rsi_2"), item.get("rsi_3"), item.get("rsi_5"), item.get("is_confirmed", 1), src, ch, item.get("holding_peak"), now_str, )) sql = self._ws_candles_upsert_sql( freeze=self._ws_candle_freeze_on_confirm(), ) # ws_candles 테이블이 존재하면 배치 INSERT (없으면 조용히 skip) self.db.conn.execute( sql, rows[0], ) if len(rows) == 1 else self._executemany_batch(rows) logger.debug("💾 [배치저장] %d봉 → DB", len(rows)) except Exception as e: logger.debug("_flush_batch 실패 (무시): %s", e) def _executemany_batch(self, rows: list) -> None: """여러 봉을 executemany 로 한 번에 INSERT.""" sql = self._ws_candles_upsert_sql( freeze=self._ws_candle_freeze_on_confirm(), ) # pymysql executemany: cursor.executemany(sql, list_of_tuples) with self.db.conn._lock: self.db.conn._ensure_connected() cur = self.db.conn._conn.cursor() cur.executemany(sql, rows) def stop(self) -> None: """기록원 스레드를 안전하게 종료합니다.""" if self._db_writer_thread and self._db_writer_thread.is_alive(): self._write_queue.put(None) # 종료 신호 self._db_writer_thread.join(timeout=5) # ------------------------------------------------------------------ # 시각 유틸 # ------------------------------------------------------------------ @staticmethod def _candle_key(tick_time_str: str, timeframe: int) -> str: """ 틱 시각과 timeframe(분)으로 봉 시작 시각 키를 반환. 반환 형식: YYYYMMDDHHMM (분 단위, 초 버림) 예) tick_time="091523", timeframe=3 → "202603020915" (09:15봉) 한투는 YYYYMMDD+HHMMSS(14자리+)로 올 수 있음 — 앞 8자리를 시분으로 읽지 말 것. 키움 FID20 는 HHMMSS(6자리). 날짜는 오늘. """ import datetime as _dt try: raw = (tick_time_str or "").strip().replace(":", "").replace("-", "").replace(" ", "") tf = int(timeframe) or 1 if len(raw) >= 14 and raw[:14].isdigit(): today = raw[:8] hh = int(raw[8:10]) mm = int(raw[10:12]) elif len(raw) >= 6 and raw[:6].isdigit(): hh = int(raw[0:2]) mm = int(raw[2:4]) today = _dt.date.today().strftime("%Y%m%d") else: return _dt.datetime.now().strftime("%Y%m%d%H%M") # timeframe 단위로 내림: 3분봉이면 09:17 → 09:15 floored_mm = (mm // tf) * tf return f"{today}{hh:02d}{floored_mm:02d}" except Exception: import datetime as _dt return _dt.datetime.now().strftime("%Y%m%d%H%M") @staticmethod def _open_bucket_ctime(tf: int, now=None) -> str: """현재 시각 기준 진행 중(미완성) 봉의 candle_time (YYYYMMDDHHMM).""" import datetime as _dt from kis_trader.engine.candle_rollup import floor_candle_time_to_tf dt0 = now if now is not None else _dt.datetime.now() raw = dt0.strftime("%Y%m%d%H%M") return floor_candle_time_to_tf(raw, int(tf) or 1) # ------------------------------------------------------------------ # RSI 계산 (확정 봉 close 리스트 기반) # ------------------------------------------------------------------ @staticmethod def _calc_rsi(closes: list, period: int) -> Optional[float]: """ close 가격 리스트에서 RSI(period) 계산. - Wilder's smoothed moving average (단순 rolling mean 사용) - 데이터 부족 시 None 반환 """ 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 _compute_rsi_set(self, closes: list) -> tuple: """RSI 2/3/5 를 한 번에 계산해서 (rsi2, rsi3, rsi5) 반환.""" return ( self._calc_rsi(closes, 2), self._calc_rsi(closes, 3), self._calc_rsi(closes, 5), ) # ------------------------------------------------------------------ # 핵심: 틱 수신 처리 # ------------------------------------------------------------------ def on_tick( self, code: str, price: float, volume: int, tick_time: str, *, market: str = "KR", timeframes: Optional[list] = None, source: str = "kis", ) -> None: """ KISWebSocketPriceCache._parse_realtime_msg 에서 틱마다 호출됨. 기본: 등록된 모든 timeframe 봉 갱신. timeframes: 해외 WS 등에서 1분만 넘길 때 (국내 호출부는 None 유지). market: KR(국내) | US(해외 HDFSCNT0) — ws_candles.market 기록용. """ if price <= 0: return mk = (market or "KR").strip().upper()[:8] or "KR" tfs = list(timeframes) if timeframes is not None else self.timeframes with self._lock: for tf in tfs: self._process_tick(code, price, volume, tick_time, tf, market=mk, source=source) def _process_tick(self, code: str, price: float, volume: int, tick_time: str, tf: int, *, market: str = "KR", source: str = "kis") -> None: """ 단일 timeframe 에 대한 틱 처리 (lock 내부에서 호출). [트랙 1] RAM 갱신만 수행, DB 호출 없음 → 블로킹 0ms [트랙 2] 봉 확정 순간에만 Queue.put_nowait() → 기록원이 비동기 배치 저장 """ key = self._bar_key(code, tf, source, "ws") new_ctime = self._candle_key(tick_time, tf) if key not in self._current: # 첫 틱: 새 봉 시작 (RAM만) inc = self._volume_is_incremental(code) self._current[key] = { "candle_time": new_ctime, "open": price, "high": price, "low": price, "close": price, "volume": int(volume) if inc else 0, "_acml_base": None if inc else int(volume), "market": market, "source": key[2], "channel": key[3], } return cur = self._current[key] cur["market"] = market or cur.get("market") or "KR" if new_ctime != cur["candle_time"]: # ── 봉 확정 (다음 틱 도착 → 정상 롤오버) ──────────────── confirmed_candle = self._confirm_current_bucket(key, cur) logger.debug( "🕯 [봉확정] %s %dM %s C=%.0f RSI3=%s", code, tf, confirmed_candle["candle_time"], confirmed_candle["close"], f"{confirmed_candle['rsi_3']:.1f}" if confirmed_candle["rsi_3"] is not None else "N/A", ) # ── 새 봉 시작 (RAM만) ────────────────────────────────── inc = self._volume_is_incremental(code) self._current[key] = { "candle_time": new_ctime, "open": price, "high": price, "low": price, "close": price, "volume": int(volume) if inc else 0, "_acml_base": None if inc else int(volume), "market": market, "source": key[2], "channel": key[3], } else: # 같은 봉: OHLCV 갱신 (RAM만, DB 쓰기 없음) cur["high"] = max(cur["high"], price) cur["low"] = min(cur["low"], price) cur["close"] = price if self._volume_is_incremental(code): # 키움 FID 15: 틱 체결량 합산 cur["volume"] = int(cur.get("volume", 0) or 0) + int(volume or 0) else: # KIS ACML_VOL: 봉 시작 대비 델타 base = int(cur.get("_acml_base", 0) or 0) cur["volume"] = max(0, int(volume or 0) - base) def _confirm_current_bucket(self, key: tuple, cur: Dict) -> Dict: """ 진행 중이던 ``_current[key]`` 봉을 확정봉으로 전환한다. (``_process_tick``의 정상 롤오버 / ``flush_stale_current_candles``의 시간경과 강제확정 양쪽에서 공통으로 사용 — 락 보유 상태에서 호출) [트랙 1] RSI 계산 + ``_confirmed`` RAM 버퍼 적재 (매수 루프 즉시 참조용) [트랙 2] DB 기록 Queue 적재 (논블로킹) 동일 ``candle_time`` 이 이미 있으면: - ``WS_CANDLE_FREEZE_ON_CONFIRM``(기본 true): OHLCV 유지(첫 확정 승), holding_peak 만 갱신 - freeze OFF: volume 더 큰 쪽 upsert (레거시) """ code = key[0] tf = key[1] source = key[2] if len(key) >= 3 else "kis" channel = key[3] if len(key) >= 4 else str(cur.get("channel") or "ws") ctime = str(cur.get("candle_time") or "")[:12] buf = self._confirmed.setdefault(key, []) closes = self._closes.setdefault(key, []) new_vol = int(cur.get("volume") or 0) idx = next( ( i for i, c in enumerate(buf) if str(c.get("candle_time") or "")[:12] == ctime ), -1, ) hp = self._holding_peak_for(code) freeze = self._ws_candle_freeze_on_confirm() if idx >= 0: old = buf[idx] old_vol = int(old.get("volume") or 0) # freeze: 이미 확정된 봉은 OHLCV 고정 (REST/재확정이 키우지 않음) if freeze or new_vol < old_vol: confirmed_candle = dict(old) confirmed_candle["is_confirmed"] = 1 confirmed_candle["source"] = str(old.get("source") or source) confirmed_candle["channel"] = str(old.get("channel") or channel) confirmed_candle["market"] = str( cur.get("market") or old.get("market") or "KR" ).upper()[:8] if hp is not None: confirmed_candle["holding_peak"] = max( float(confirmed_candle.get("holding_peak") or 0), float(hp), ) # freeze 시 high 도 동결 — peak 만 메타로 보관 if not freeze: confirmed_candle["high"] = max( float(confirmed_candle.get("high") or 0), float(hp), ) buf[idx] = confirmed_candle return confirmed_candle low_cands = [ x for x in (float(old.get("low") or 0), float(cur["low"])) if x > 0 ] confirmed_candle = { "code": code, "market": str(cur.get("market") or old.get("market") or "KR").upper()[:8], "tf": tf, "candle_time": ctime, "open": float(old.get("open") or cur["open"]), "high": max(float(old.get("high") or 0), float(cur["high"])), "low": min(low_cands) if low_cands else float(cur["low"]), "close": float(cur["close"]), "volume": max(old_vol, new_vol), "is_confirmed": 1, "source": str(cur.get("source") or old.get("source") or source), "channel": str(cur.get("channel") or old.get("channel") or channel), } if hp is not None: confirmed_candle["holding_peak"] = hp confirmed_candle["high"] = max( float(confirmed_candle["high"]), float(hp), ) buf[idx] = confirmed_candle closes[:] = [float(c.get("close") or 0) for c in buf] rsi2, rsi3, rsi5 = self._compute_rsi_set(closes[: idx + 1]) confirmed_candle["rsi_2"] = rsi2 confirmed_candle["rsi_3"] = rsi3 confirmed_candle["rsi_5"] = rsi5 buf[idx] = confirmed_candle else: closes.append(float(cur["close"])) if len(closes) > self._ram_buffer_max: closes.pop(0) rsi2, rsi3, rsi5 = self._compute_rsi_set(closes) confirmed_candle = { "code": code, "market": str(cur.get("market") or "KR").upper()[:8], "tf": tf, "candle_time": ctime, "open": cur["open"], "high": cur["high"], "low": cur["low"], "close": cur["close"], "volume": new_vol, "rsi_2": rsi2, "rsi_3": rsi3, "rsi_5": rsi5, "is_confirmed": 1, "source": cur.get("source", source), "channel": cur.get("channel", channel), } if hp is not None: confirmed_candle["holding_peak"] = hp confirmed_candle["high"] = max(cur["high"], hp) buf.append(confirmed_candle) if len(buf) > self._ram_buffer_max: buf.pop(0) closes[:] = [float(c.get("close") or 0) for c in buf] try: self._write_queue.put_nowait(confirmed_candle) except queue.Full: logger.warning("⚠️ CandleAggregator 쓰기 Queue 가득참 — 봉 1개 DROP (코드: %s)", code) return confirmed_candle # ------------------------------------------------------------------ # 저유동 종목 봉 강제확정: 다음 체결 틱이 안 와도 봉주기 경과 시 확정 # ------------------------------------------------------------------ def flush_stale_current_candles(self) -> int: """ ``_current``(진행 중 봉)가 자기 봉주기(tf분)를 이미 지났는데도 다음 체결 틱이 없어 확정되지 못한 채 머물러 있으면, 그 시점까지 쌓인 데이터로 강제 확정한다. 배경: 기존엔 "다음 틱이 와야 봉 확정"이라, 저유동 종목이 한동안 무거래면 이미 끝난 봉도 다음 체결 전까지 무한정 미확정 상태로 남아 매수 신호 인식이 그만큼 밀렸다(2026-07-08 원티드랩 14분 지연 사례 — 3분봉인데 다음 틱이 14분 뒤에야 와서 신호가 14분 늦게 잡힘). 이 함수는 heartbeat 성격으로 주기 호출되어 그 지연을 봉주기(3분) 이내로 되돌린다. 확정되는 값 자체는 실제 체결 데이터 그대로라 백테(ws_candles 그대로 재생)와 내용 차이는 없다 — "언제 인지하냐"만 앞당긴다. ``WS_CANDLE_FORCE_CONFIRM_ENABLED=false`` 로 즉시 롤백 가능. """ if not get_env_bool("WS_CANDLE_FORCE_CONFIRM_ENABLED", True): return 0 grace_sec = get_env_float("WS_CANDLE_FORCE_CONFIRM_GRACE_SEC", 0.0) import datetime as _dt now = _dt.datetime.now() confirmed_n = 0 with self._lock: for key in list(self._current.keys()): cur = self._current.get(key) if not cur: continue code = key[0] tf = key[1] try: bucket_start = _dt.datetime.strptime(cur["candle_time"], "%Y%m%d%H%M") except Exception: continue bucket_end = bucket_start + _dt.timedelta(minutes=tf) if now < bucket_end + _dt.timedelta(seconds=grace_sec): continue confirmed_candle = self._confirm_current_bucket(key, cur) self._current.pop(key, None) confirmed_n += 1 logger.info( "⏱ [봉강제확정] %s %dM %s C=%.0f — 다음 체결 없음(봉주기 경과) → 즉시 확정", code, tf, confirmed_candle["candle_time"], confirmed_candle["close"], ) return confirmed_n # ------------------------------------------------------------------ # 재접속 갭 보정: REST get_minute_chart 로 빈 봉 채우기 # ------------------------------------------------------------------ def fill_gap_from_rest(self, code: str, tf: int, rest_df) -> int: """ WS 재접속 후 빠진 봉 구간을 REST 분봉 데이터로 채움. [트랙 1] close 가격을 _closes / _confirmed 에 넣어 RSI 웜업 [트랙 2] 봉 dict를 Queue 에 넣어 기록원이 배치로 DB 저장 (source=kiwoom, channel=rest) Args: code : 종목코드 tf : timeframe (분) rest_df : get_minute_chart 반환 DataFrame (오래된→최신 순 정렬 필요) 컬럼: time, open, high, low, close, volume Returns: 채워진 봉 수 """ if rest_df is None or rest_df.empty: return 0 # 진행 중 분봉은 confirmed 에 넣지 않음 — merge_confirmed_bars 에서도 재필터. skip_incomplete = get_env_bool("WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET", True) open_bucket = self._open_bucket_ctime(tf) if skip_incomplete else "" rows: list = [] skipped_open = 0 for _, row in rest_df.iterrows(): ctime = str(row.get("time", ""))[:12] if not ctime or len(ctime) < 12: continue if open_bucket and ctime >= open_bucket: skipped_open += 1 continue close = float(row.get("close", 0) or 0) if close <= 0: continue rows.append({ "candle_time": ctime, "open": float(row.get("open", close) or close), "high": float(row.get("high", close) or close), "low": float(row.get("low", close) or close), "close": close, "volume": int(float(row.get("volume", 0) or 0)), "source": "kiwoom", "channel": "rest", }) if skipped_open: logger.info( "⏭ [갭보정] %s %dM 진행분(>=%s) %d봉 confirmed 제외 (매수 직전봉 왜곡 방지)", code, tf, open_bucket, skipped_open, ) return self.merge_confirmed_bars(code, tf, rows, log_tag="REST") def merge_confirmed_bars( self, code: str, tf: int, bars: list, *, log_tag: str = "merge", skip_incomplete_bucket: Optional[bool] = None, now=None, ) -> int: """ 확정봉 리스트를 RAM(+DB 큐)에 병합. - 신규 candle_time → insert - 기존 candle_time → ``WS_CANDLE_FREEZE_ON_CONFIRM``(기본 true) 이면 skip (freeze OFF: volume 더 클 때만 OHLCV upsert — 레거시) - ``WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET``(기본 true): 진행 중 버킷 (``candle_time >= 현재 봉시작``) 은 insert/update 하지 않고, 이미 RAM 에 있으면 제거. 장초 REST 미완성봉 → 직전봉% 몸통 왜곡 방지. """ if not bars: # bars 비어도 진행분 purge 는 수행 (재갭보정 정리) pass if skip_incomplete_bucket is None: skip_incomplete_bucket = get_env_bool( "WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET", True, ) freeze = self._ws_candle_freeze_on_confirm() open_bucket = ( self._open_bucket_ctime(tf, now=now) if skip_incomplete_bucket else "" ) from kis_trader.ws.candle_series import normalize_source_channel first_src, first_ch = "kiwoom", "rest" if bars: first_src, first_ch = normalize_source_channel( str(bars[0].get("source") or ""), str(bars[0].get("channel") or ""), ) # lock 밖에서 DB 조회 (재시작 후 RAM 공백 → REST 가 DB 확정봉을 덮는 것 방지) db_frozen: Dict[str, Dict] = {} if freeze and bars and self.db is not None: want_times = [] for row in bars: ct = str(row.get("candle_time") or row.get("time") or "")[:12] if ct and len(ct) >= 12: want_times.append(ct) db_frozen = self._load_confirmed_ohlcv_from_db( code, tf, want_times, source=first_src, channel=first_ch, ) with self._lock: key = self._bar_key(code, tf, first_src, first_ch) closes = self._closes.setdefault(key, []) conf_buf = self._confirmed.setdefault(key, []) # 이미 들어온 진행분 confirmed 제거 (갭보정 직후·롤업 공통) purged = 0 if open_bucket and conf_buf: kept = [] for c in conf_buf: ct = str(c.get("candle_time") or "")[:12] if ct and ct >= open_bucket: purged += 1 continue kept.append(c) if purged: conf_buf[:] = kept closes[:] = [float(c.get("close") or 0) for c in conf_buf] logger.info( "🧹 [갭보정] %s %dM 진행분 confirmed %d봉 제거 (>=%s)", code, tf, purged, open_bucket, ) if not bars: return 0 by_time = { str(c.get("candle_time", ""))[:12]: i for i, c in enumerate(conf_buf) if str(c.get("candle_time", ""))[:12] } inserted = 0 updated = 0 skipped_freeze = 0 seeded_db = 0 skipped_open = 0 for row in bars: ctime = str(row.get("candle_time") or row.get("time") or "")[:12] if not ctime or len(ctime) < 12: continue if open_bucket and ctime >= open_bucket: skipped_open += 1 continue close = float(row.get("close", 0) or 0) if close <= 0: continue new_vol = int(float(row.get("volume", 0) or 0)) src, ch = normalize_source_channel( str(row.get("source") or first_src), str(row.get("channel") or first_ch), ) idx = by_time.get(ctime) if idx is not None: # freeze-on-confirm: 첫 확정본 유지 (갭보정/백필이 lookback 시리즈를 키우지 않음) if freeze: skipped_freeze += 1 continue old = conf_buf[idx] old_vol = int(old.get("volume") or 0) if new_vol <= old_vol: continue low_cands = [ x for x in ( float(old.get("low") or 0), float(row.get("low", close) or close), ) if x > 0 ] candle = { "code": code, "tf": tf, "candle_time": ctime, "open": float(old.get("open") or row.get("open", close) or close), "high": max( float(old.get("high") or 0), float(row.get("high", close) or close), ), "low": min(low_cands) if low_cands else float( row.get("low", close) or close ), "close": close, "volume": new_vol, "is_confirmed": 1, "source": src, "channel": ch, } if old.get("holding_peak") is not None: candle["holding_peak"] = old.get("holding_peak") conf_buf[idx] = candle closes[:] = [float(c.get("close") or 0) for c in conf_buf] rsi2, rsi3, rsi5 = self._compute_rsi_set(closes[: idx + 1]) candle["rsi_2"] = rsi2 candle["rsi_3"] = rsi3 candle["rsi_5"] = rsi5 conf_buf[idx] = candle try: self._write_queue.put_nowait(candle) except queue.Full: pass updated += 1 continue # RAM 에 없고 DB 에 확정봉이 있으면 → DB 값으로만 시드 (REST로 덮지 않음) if freeze and ctime in db_frozen: dr = db_frozen[ctime] d_close = float(dr.get("close") or 0) if d_close <= 0: d_close = close closes.append(d_close) if len(closes) > self._ram_buffer_max: closes.pop(0) rsi2 = dr.get("rsi_2") rsi3 = dr.get("rsi_3") rsi5 = dr.get("rsi_5") if rsi2 is None and rsi3 is None and rsi5 is None: rsi2, rsi3, rsi5 = self._compute_rsi_set(closes) candle = { "code": code, "tf": tf, "candle_time": ctime, "open": float(dr.get("open") or d_close), "high": float(dr.get("high") or d_close), "low": float(dr.get("low") or d_close), "close": d_close, "volume": int(float(dr.get("volume") or 0)), "rsi_2": rsi2, "rsi_3": rsi3, "rsi_5": rsi5, "is_confirmed": 1, "source": str(dr.get("source") or first_src)[:16], "channel": str(dr.get("channel") or first_ch)[:10], } if dr.get("holding_peak") is not None: candle["holding_peak"] = dr.get("holding_peak") conf_buf.append(candle) by_time[ctime] = len(conf_buf) - 1 if len(conf_buf) > self._ram_buffer_max: conf_buf.sort(key=lambda x: str(x.get("candle_time", ""))) while len(conf_buf) > self._ram_buffer_max: conf_buf.pop(0) closes[:] = [float(c["close"]) for c in conf_buf] by_time = { str(c.get("candle_time", ""))[:12]: i for i, c in enumerate(conf_buf) if str(c.get("candle_time", ""))[:12] } # DB 이미 있음 → 쓰기 큐 불필요 seeded_db += 1 continue closes.append(close) if len(closes) > self._ram_buffer_max: closes.pop(0) rsi2, rsi3, rsi5 = self._compute_rsi_set(closes) candle = { "code": code, "tf": tf, "candle_time": ctime, "open": float(row.get("open", close) or close), "high": float(row.get("high", close) or close), "low": float(row.get("low", close) or close), "close": close, "volume": new_vol, "rsi_2": rsi2, "rsi_3": rsi3, "rsi_5": rsi5, "is_confirmed": 1, "source": src, "channel": ch, } conf_buf.append(candle) by_time[ctime] = len(conf_buf) - 1 if len(conf_buf) > self._ram_buffer_max: conf_buf.sort(key=lambda x: str(x.get("candle_time", ""))) while len(conf_buf) > self._ram_buffer_max: conf_buf.pop(0) closes[:] = [float(c["close"]) for c in conf_buf] by_time = { str(c.get("candle_time", ""))[:12]: i for i, c in enumerate(conf_buf) if str(c.get("candle_time", ""))[:12] } try: self._write_queue.put_nowait(candle) except queue.Full: pass inserted += 1 if inserted or updated or seeded_db: conf_buf.sort(key=lambda x: str(x.get("candle_time", ""))) closes[:] = [float(c["close"]) for c in conf_buf] if skipped_open and log_tag: logger.debug( "⏭ [갭보정] %s %dM %s 진행분 %d봉 skip (>=%s)", code, tf, log_tag, skipped_open, open_bucket, ) if inserted or updated or skipped_freeze or seeded_db: logger.info( "🔧 [갭보정] %s %dM → %s insert=%d update=%d freeze_skip=%d db_seed=%d RAM+DB큐", code, tf, log_tag, inserted, updated, skipped_freeze, seeded_db, ) return inserted + updated + seeded_db @staticmethod def _live_candle_source_order() -> tuple: """봉 읽기 순서 — 메인 WS → 2차 WS → LS WS → kiwoom REST → 롤업. CANDLE_SOURCE=kis|kiwoom 이면 그 메인, 빈값이면 LIVE_TICK_PROVIDER. """ from kis_trader.ws.candle_series import live_read_pairs return live_read_pairs() def rollup_tf_from_1m(self, code: str, target_tf: int = 3) -> int: """ RAM 1분 확정봉 → target_tf 분봉 재합성 후 구멍만 보강. 꼬리(3M) 트리거 웜업: 1M REST 1회만으로 3M 준비 (WS_GAP_ROLLUP_3M_FROM_1M). 1M 소스 우선순위는 실매 읽기와 동일(메인 WS → 키움 REST). 저장은 source=kiwoom, channel=rollup (증권사 WS 행과 분리). """ from kis_trader.engine.candle_rollup import rollup_1m_bars_to_tf tf = int(target_tf) if tf <= 1: return 0 pairs = tuple( p for p in self._live_candle_source_order() if p[1] != "rollup" ) with self._lock: seen: dict = {} for src, ch in pairs: for bar in self._confirmed.get(self._bar_key(code, 1, src, ch), []): ct = str(bar.get("candle_time") or "")[:12] if ct and ct not in seen: seen[ct] = bar bars_1m = sorted(seen.values(), key=lambda x: str(x.get("candle_time") or "")) if not bars_1m: return 0 rolled = rollup_1m_bars_to_tf(bars_1m, tf) for b in rolled or []: b["source"] = "kiwoom" b["channel"] = "rollup" return self.merge_confirmed_bars( code, tf, rolled, log_tag=f"rollup_1m→{tf}M", ) # ------------------------------------------------------------------ # [트랙 1] RAM 버퍼 조회 — 매수/매도 루프에서 직접 호출 (DB 조회 없음) # ------------------------------------------------------------------ # NOTE: source="" (기본) → LIVE_TICK_PROVIDER 우선으로 소스 병합 # 재시작 직후 REST·WS 갭보정 데이터를 전략/롤업이 즉시 인식하기 위함. # source를 명시할 경우 (예: source="kis") 해당 소스만 조회. def _copy_raw_for_merge(self, code: str, tf: int) -> list: """락 안에서 호출 — 우선순위 순으로 봉을 이어붙임 (정렬 없음). 락 보유 시간을 최소화하기 위해 정렬은 락 밖(_merge_from_raw)에서 수행. 같은 candle_time이면 앞에 나온 소스(메인)가 이김. """ items: list = [] for src, ch in self._live_candle_source_order(): if src == "ls" and ch == "ws": fn = getattr(self, "_ls_bars_fn", None) if callable(fn): try: extra = fn(code, tf) or [] except Exception: extra = [] for b in extra: row = dict(b) row["source"] = "ls" row["channel"] = "ws" items.append(row) continue buf = self._confirmed.get(self._bar_key(code, tf, src, ch)) if buf: items.extend(buf) return items def _merge_from_raw(self, raw_items: list, code: str = "", tf: int = 1) -> list: """락 밖에서 호출 — 쓰레기 WS 봉 skip 후 시간순.""" from kis_trader.ws.candle_series import dedupe_by_read_pairs, live_read_pairs ticks = None live = False lookup = getattr(self, "_tick_lookup", None) if callable(lookup) and code: try: ticks = lookup(code) or [] live = True except Exception: ticks = None return dedupe_by_read_pairs( raw_items, live_read_pairs(), ticks=ticks, tf_min=int(tf or 1), missing_policy="hole", live_cover=bool(live), ) def _merge_confirmed_all_sources(self, code: str, tf: int) -> list: """모든 소스의 확정봉을 candle_time 기준으로 병합·정렬. NOTE: 락 안에서 호출하는 레거시 경로. 신규 코드는 _copy_raw_for_merge + _merge_from_raw 사용. """ return self._merge_from_raw(self._copy_raw_for_merge(code, tf), code, tf) def get_latest_confirmed(self, code: str, tf: int, source: str = "") -> Optional[dict]: """ 가장 최근 확정된 봉(완성된 마지막 봉)을 반환. source="" (기본) → 모든 소스 중 가장 최근 봉. None이면 아직 봉이 확정되지 않음 (장 초반 등). """ if source: with self._lock: buf = self._confirmed.get(self._bar_key(code, tf, source)) return buf[-1] if buf else None with self._lock: raw = self._copy_raw_for_merge(code, tf) merged = self._merge_from_raw(raw, code, tf) return merged[-1] if merged else None def get_prev_confirmed(self, code: str, tf: int, source: str = "") -> Optional[dict]: """직전 확정봉 (최신에서 2번째). 패턴 확인용 (현재봉 - 1).""" if source: with self._lock: buf = self._confirmed.get(self._bar_key(code, tf, source)) return buf[-2] if buf and len(buf) >= 2 else None with self._lock: raw = self._copy_raw_for_merge(code, tf) merged = self._merge_from_raw(raw, code, tf) return merged[-2] if len(merged) >= 2 else None def get_candles(self, code: str, tf: int, n: int = 10, source: str = "") -> list: """최근 n개 확정 봉 리스트 반환 (오래된→최신 순). source="" (기본) → 모든 소스 병합 후 최근 n개. 락 안에서는 list 복사만, 정렬은 락 밖에서 수행 (블로킹 최소화). """ if source: with self._lock: buf = self._confirmed.get(self._bar_key(code, tf, source), []) return list(buf[-n:]) with self._lock: raw = self._copy_raw_for_merge(code, tf) merged = self._merge_from_raw(raw, code, tf) return merged[-n:] def get_confirmed_count(self, code: str, tf: int, source: str = "") -> int: """확정된 봉 수 (RSI 안정화 여부 확인용). source="" (기본) → 모든 소스 병합 후 고유 봉 수. """ if source: with self._lock: return len(self._confirmed.get(self._bar_key(code, tf, source), [])) with self._lock: raw = self._copy_raw_for_merge(code, tf) return len(self._merge_from_raw(raw, code, tf)) def get_current_candle(self, code: str, tf: int, source: str = "kis") -> Optional[dict]: """ 현재 진행 중인 봉(미확정, is_confirmed=0) 반환. RSI는 포함되지 않음 (확정 봉 기준으로만 계산). """ with self._lock: return dict(self._current.get(self._bar_key(code, tf, source), {})) or None def get_rsi(self, code: str, tf: int, period: int = 3, source: str = "kis") -> Optional[float]: """ 최신 확정 봉의 RSI(period) 값 반환. period: 2, 3, 5 중 하나 (스캘핑 단타용 초단기 RSI) """ candle = self.get_latest_confirmed(code, tf, source=source) if candle is None: return None return candle.get(f"rsi_{period}") def remove_code(self, code: str) -> None: """ 유니버스에서 빠진 종목의 RAM 버퍼를 정리합니다. _sync_subscriptions()에서 구독 해제 시 같이 호출. 등록된 모든 timeframe 의 _confirmed / _closes / _current 를 삭제. """ with self._lock: for store in (self._confirmed, self._closes, self._current): for key in list(store.keys()): if key and key[0] == code: store.pop(key, None) logger.debug("🗑️ CandleAggregator RAM 정리: %s", code) # ====================================================================== # 키움 REST API 공통 유틸 — 봇 시작 시 WS 갭보정용 # KIS get_minute_chart() 는 당일봉만 제공하지만, # 키움 ka10080 은 1회 호출에 최대 900봉(≈6개월) + 페이지네이션 지원. # ====================================================================== # ────────────────────────────────────────────────────────────────────── # 키움 토큰 매니저 싱글톤 풀 # # kiwoom_rest_api/auth/token.py 의 TokenManager 와 동일한 설계: # - _is_access_token_valid(): expires_dt 기반 정확한 만료 판별 (30s 버퍼) # - _request_new_token(): 만료 시에만 발급 (au10001 rate limit 방지) # # 차이점: DB에서 읽은 app_key/secret 직접 주입 (환경변수 불필요) # 싱글톤 풀: 캐시 키(domain:appkey앞8자) → 인스턴스 재사용 # 종목×timeframe마다 새 인스턴스를 만들지 않음 # ────────────────────────────────────────────────────────────────────── import threading as _threading from datetime import datetime as _datetime, timedelta as _timedelta _kiwoom_token_lock = _threading.Lock() _kiwoom_managers: Dict[str, "KiwoomTokenManager"] = {} # 싱글톤 풀 # ka10080 전역 동시호출 상한 — 키움 유량=5. 워커 N개여도 이 세마포어로 직렬화. _ka10080_sem: Optional[_threading.Semaphore] = None _ka10080_sem_lock = _threading.Lock() _ka10080_sem_n: int = 0 def _ka10080_semaphore() -> _threading.Semaphore: """키움 ka10080 동시 요청 수 제한 (유량=5 미만으로 유지).""" global _ka10080_sem, _ka10080_sem_n n = max(1, min(int(get_env_int("KIWOOM_KA10080_MAX_INFLIGHT", 2)), 4)) with _ka10080_sem_lock: if _ka10080_sem is None or _ka10080_sem_n != n: _ka10080_sem = _threading.Semaphore(n) _ka10080_sem_n = n return _ka10080_sem def _ka10080_sem_timeout_sec(*, purpose: str = "") -> float: """ka10080/ka10007 공유 sem 대기 상한(초). 0=무제한(레거시).""" key = "KIWOOM_KA10007_SEM_TIMEOUT_SEC" if purpose == "ka10007" else "KIWOOM_KA10080_SEM_TIMEOUT_SEC" default = 2.0 if purpose == "ka10007" else 12.0 try: return max(0.0, float(get_env_float(key, default) or 0.0)) except Exception: return default def _ka10080_acquire( sem: _threading.Semaphore, *, purpose: str, code: str = "", ) -> bool: """갭보정(ka10080)·매도 REST(ka10007) sem — 무한 대기 금지.""" wait = _ka10080_sem_timeout_sec(purpose=purpose) if wait <= 0: sem.acquire() return True if sem.acquire(timeout=wait): return True logger.warning( "⚠️ 키움 REST sem 타임아웃 purpose=%s code=%s wait=%.0fs " "(갭 ka10080 점유 중 — 매도·매수 루프 기아 방지)", purpose, code or "-", wait, ) return False def _ka10080_is_rate_limit(http_st: int, rc: Any, msg: str) -> bool: if int(http_st or 0) == 429: return True try: if int(rc) == 5: return True except Exception: pass m = str(msg or "") return ("1700" in m) or ("허용된 요청" in m) or ("유량" in m) class KiwoomTokenManager: """ kiwoom_rest_api/auth/token.py TokenManager 와 동일한 로직 — DB 키 직접 주입 버전. - get_token(): 유효한 토큰이면 바로 반환, 만료 시에만 재발급 (au10001 방지) - expires_dt 를 정확히 파싱 (하드코딩 23h 아님) - 30초 버퍼: 경계 케이스 방지 (TokenManager 와 동일) - 스레드 안전: 인스턴스당 Lock 보유 """ 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 self._domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com" self._mode_str = "모의" if is_mock else "실전" self._token: Optional[str] = None self._expiry: Optional[_datetime] = None self._lock = _threading.Lock() def _is_valid(self) -> bool: """만료 30초 전까지 유효 (TokenManager._is_access_token_valid 와 동일)""" if not self._token or not self._expiry: return False return _datetime.now() < self._expiry - _timedelta(seconds=30) def _request_new_token(self) -> None: """토큰 신규 발급 — 만료 시에만 호출됨""" resp = requests.post( f"https://{self._domain}/oauth2/token", json={ "grant_type": "client_credentials", "appkey": self._app_key, "secretkey": self._app_secret, }, timeout=10, ) data = resp.json() token = (data.get("token") or data.get("access_token") or "").strip() if not token: raise RuntimeError(f"키움 토큰 발급 실패 [{self._mode_str}]: {data}") # expires_dt 정확히 파싱 (TokenManager._update_token_info 와 동일 로직) exp_s = data.get("expires_dt", "") try: self._expiry = _datetime.strptime(str(exp_s), "%Y%m%d%H%M%S") except Exception: # expires_in 도 없으면 24h 기본 (보수적 fallback) exp_in = data.get("expires_in", 86400) self._expiry = _datetime.now() + _timedelta(seconds=int(exp_in)) self._token = token logger.info( "✅ 키움 토큰 발급 완료 [%s] (앞8자: %s…, 만료: %s)", self._mode_str, token[:8], self._expiry.strftime("%Y-%m-%d %H:%M:%S"), ) def invalidate(self, reason: str = "") -> None: """서버가 토큰을 거부했을 때 로컬 캐시를 비운다 (만료 전 재사용 금지). WS LOGIN rc=805004 등 — expires_dt 전이라도 서버 측 무효면 같은 토큰을 재전송하면 재연결만 소모하고 영구 비활성으로 간다. """ with self._lock: self._token = None self._expiry = None if reason: logger.warning( "♻️ 키움 토큰 캐시 무효화 [%s]: %s", self._mode_str, reason, ) def get_token(self) -> Optional[str]: """ 유효한 토큰 반환. 만료 시에만 재발급. kiwoom_rest_api TokenManager.get_token() 호환 인터페이스. """ with self._lock: if not self._is_valid(): try: self._request_new_token() except Exception as e: logger.warning("⚠️ 키움 토큰 발급 예외 [%s]: %s", self._mode_str, e) return None return self._token def _kiwoom_token_cache_key(kiwoom_key: str, is_mock: bool) -> str: domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com" return f"{domain}:{str(kiwoom_key or '')[:8]}" def _get_kiwoom_token_cached( kiwoom_key: str, kiwoom_secret: str, is_mock: bool, ) -> Optional[str]: """ KiwoomTokenManager 싱글톤 풀에서 인스턴스를 꺼내 토큰 반환. 동일 키 조합은 같은 인스턴스를 재사용 → au10001 rate limit 방지. """ cache_key = _kiwoom_token_cache_key(kiwoom_key, is_mock) with _kiwoom_token_lock: if cache_key not in _kiwoom_managers: _kiwoom_managers[cache_key] = KiwoomTokenManager( kiwoom_key, kiwoom_secret, is_mock=is_mock ) mgr = _kiwoom_managers[cache_key] return mgr.get_token() def _invalidate_kiwoom_token_cached( kiwoom_key: str, kiwoom_secret: str, is_mock: bool, reason: str = "", ) -> None: """LOGIN/REST 토큰 거부 시 싱글톤 캐시 무효화 — 다음 get_token 이 재발급.""" cache_key = _kiwoom_token_cache_key(kiwoom_key, is_mock) with _kiwoom_token_lock: mgr = _kiwoom_managers.get(cache_key) if mgr is not None: mgr.invalidate(reason or "token rejected") def _get_kiwoom_creds(db) -> tuple: """ env_auth_config 최신 행에서 키움 앱키/시크릿 반환. 시세 REST(갭보정 ka10080 · 매도 4차 ka10007)는 **매매 KIS_MOCK 과 분리**. 기본 ``KIWOOM_WS_FORCE_REAL=true`` → 실키·api.kiwoom.com (키움 WS·WSManager._get_kiwoom_credentials 와 동일). Returns: (app_key, app_secret, is_mock) — 키 없으면 (None, None, False) """ try: r = db.get_auth_env() if hasattr(db, "get_auth_env") else {} if not r: return None, None, False force_real_str = ( get_env_from_db("KIWOOM_WS_FORCE_REAL", "true") or "true" ).strip().lower() force_real = force_real_str in ("true", "1", "yes", "y", "on") if force_real: is_mock = False else: kw_mock_raw = str(r.get("KIWOOM_MOCK") or get_env_from_db("KIWOOM_MOCK", "") or "").strip().lower() if kw_mock_raw in ("true", "1", "yes", "y", "on"): is_mock = True elif kw_mock_raw in ("false", "0", "no", "n", "off"): is_mock = False else: is_mock = str(r.get("KIS_MOCK", "true")).lower() in ("true", "1", "yes") if is_mock: key = str(r.get("KIWOOM_APP_KEY_MOCK", "") or "").strip() secret = str(r.get("KIWOOM_APP_SECRET_MOCK", "") or "").strip() else: key = str(r.get("KIWOOM_APP_KEY_REAL", "") or "").strip() secret = str(r.get("KIWOOM_APP_SECRET_REAL", "") or "").strip() # 레거시 필드 폴백 (KIWOOM_APP_KEY) if not key or not secret: key = str(r.get("KIWOOM_APP_KEY", "") or "").strip() secret = str(r.get("KIWOOM_APP_SECRET", "") or "").strip() if not key or not secret: return None, None, is_mock return key, secret, is_mock except Exception as e: logger.debug("키움 크레덴셜 조회 실패: %s", e) return None, None, False def get_kiwoom_candles_df( code: str, tf_min: int, kiwoom_key: str, kiwoom_secret: str, is_mock: bool = False, n: int = 120, ) -> "object": # pd.DataFrame """ 키움 REST API (ka10080 — 주식분봉차트조회) 로 분봉 조회. fill_gap_from_rest() 호환 DataFrame 반환. ★ 유량 한도: 키움 ka10080 ``유량=5`` (return_code=5 / msg 1700). 예전: rc 무시 + records=[] → "REST 빈 응답" 오인 → 갭보정 force 재큐 폭주. 지금: 전역 세마포어(기본 동시 2) + rc=5 백오프 재시도. Returns: pd.DataFrame — time/open/high/low/close/volume, 오래된→최신 순 빈 DataFrame (오류·진짜 무데이터) """ try: import pandas as pd except ImportError: logger.error("pandas 미설치 → 키움 갭보정 불가") return None # type: ignore[return-value] domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com" token = _get_kiwoom_token_cached(kiwoom_key, kiwoom_secret, is_mock) if not token: return pd.DataFrame() base_url = f"https://{domain}/api/dostk/chart" headers = { "content-type": "application/json;charset=UTF-8", "appkey": kiwoom_key, "appsecret": kiwoom_secret, "authorization": f"Bearer {token}", "api-id": "ka10080", "cont-yn": "N", "next-key": "", } body = { "stk_cd": code, "tic_scope": str(tf_min), "upd_stkpc_tp": "1", } rate_retries = max(1, int(get_env_int("KIWOOM_KA10080_RATE_RETRIES", 5))) rate_base = float(get_env_float("KIWOOM_KA10080_RATE_RETRY_BASE_SEC", 1.2)) page_sleep = float(get_env_float("KIWOOM_KA10080_PAGE_SLEEP_SEC", 0.35)) sem = _ka10080_semaphore() rows: list = [] while len(rows) < n: data = None resp = None for attempt in range(rate_retries): if not _ka10080_acquire(sem, purpose="ka10080", code=code): return pd.DataFrame() try: resp = requests.post(base_url, json=body, headers=headers, timeout=15) http_st = int(resp.status_code) try: data = resp.json() if resp.content else {} except Exception: data = {} except Exception as e: logger.warning("⚠️ 키움 ka10080 조회 실패 (%s %dM): %s", code, tf_min, e) return pd.DataFrame() finally: sem.release() rc_raw = (data or {}).get("return_code", 0) msg = str((data or {}).get("return_msg") or "").strip() if _ka10080_is_rate_limit(http_st, rc_raw, msg): if attempt + 1 < rate_retries: wait = rate_base * (attempt + 1) + random.uniform(0.2, 0.8) logger.warning( "⚠️ 키움 ka10080 유량초과 (%s %dM) rc=%s — %.1fs 후 재시도 %d/%d " "(유량=5 · 빈응답 오인 금지)", code, tf_min, rc_raw, wait, attempt + 2, rate_retries, ) time.sleep(wait) continue logger.error( "🚨 키움 ka10080 유량초과 재시도 소진 (%s %dM) rc=%s msg=%s", code, tf_min, rc_raw, msg[:120], ) return pd.DataFrame() try: rc_i = int(rc_raw) except Exception: rc_i = 0 if rc_i != 0: logger.warning( "⚠️ 키움 ka10080 API오류 (%s %dM) rc=%s msg=%s", code, tf_min, rc_raw, msg[:120], ) return pd.DataFrame() break # 성공(rc=0) — 페이지 처리 else: return pd.DataFrame() records = (data or {}).get("stk_min_pole_chart_qry") or [] if not records: logger.info( "키움 ka10080 반환 행 없음 (%s %dM) rc=0 msg=%s", code, tf_min, str((data or {}).get("return_msg") or "")[:80], ) break for rec in records: raw_dt = str(rec.get("cntr_tm", "")) if len(raw_dt) < 12: continue try: rows.append({ "time": raw_dt[:12], "open": abs(float(rec.get("open_pric", 0) or 0)), "high": abs(float(rec.get("high_pric", 0) or 0)), "low": abs(float(rec.get("low_pric", 0) or 0)), "close": abs(float(rec.get("cur_prc", 0) or 0)), "volume": abs(int(float(rec.get("trde_qty", 0) or 0))), }) except (ValueError, TypeError): continue if len(rows) >= n: break cont_yn = str((resp.headers if resp is not None else {}).get("cont-yn", "N")).upper() next_key = str((resp.headers if resp is not None else {}).get("next-key", "")).strip() if cont_yn != "Y" or not next_key or len(rows) >= n: break headers["cont-yn"] = cont_yn headers["next-key"] = next_key if page_sleep > 0: time.sleep(page_sleep) if not rows: return pd.DataFrame() df = pd.DataFrame(rows[:n][::-1]) logger.info("✅ 키움 %dM 갭보정 데이터: %s %d봉", tf_min, code, len(df)) return df def fetch_kiwoom_cur_prc_ka10007( code: str, kiwoom_key: str, kiwoom_secret: str, is_mock: bool = False, ) -> float: """키움 ka10007(시세표성정보) 현재가. 4차 REST 전용. ka10001(유통주식 워커)와 TR 분리. ka10080 과 같은 세마포어로 유량 공유. 핫패스이므로 재시도는 짧게. 실패=0. """ code = str(code or "").strip() if not code or not kiwoom_key or not kiwoom_secret: return 0.0 token = _get_kiwoom_token_cached(kiwoom_key, kiwoom_secret, is_mock) if not token: return 0.0 domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com" url = f"https://{domain}/api/dostk/mrkcond" headers = { "content-type": "application/json;charset=UTF-8", "appkey": kiwoom_key, "appsecret": kiwoom_secret, "authorization": f"Bearer {token}", "api-id": "ka10007", "cont-yn": "N", "next-key": "", } body = {"stk_cd": code} sem = _ka10080_semaphore() data: Dict[str, Any] = {} http_st = 0 if not _ka10080_acquire(sem, purpose="ka10007", code=code): return 0.0 try: resp = requests.post(url, json=body, headers=headers, timeout=8) http_st = int(resp.status_code) try: data = resp.json() if resp.content else {} except Exception: data = {} except Exception as e: logger.warning("키움 ka10007 조회 실패 %s: %s", code, e) return 0.0 finally: sem.release() rc_raw = (data or {}).get("return_code", 0) msg = str((data or {}).get("return_msg") or "").strip() if _ka10080_is_rate_limit(http_st, rc_raw, msg): logger.warning( "키움 ka10007 유량 %s http=%s rc=%s msg=%s — 재시도 안 함", code, http_st, rc_raw, msg[:80], ) return 0.0 try: if int(rc_raw) != 0: logger.debug("키움 ka10007 API오류 %s rc=%s %s", code, rc_raw, msg[:80]) return 0.0 except Exception: return 0.0 raw = data.get("cur_prc") if isinstance(data, dict) else None if raw is None and isinstance(data, dict): inner = data.get("output") or data.get("data") or {} if isinstance(inner, dict): raw = inner.get("cur_prc") try: return abs(float(str(raw or "0").replace(",", ""))) except (TypeError, ValueError): return 0.0 def fetch_kiwoom_stock_meta_detail( code: str, kiwoom_key: str, kiwoom_secret: str, is_mock: bool = False, *, max_retries: Optional[int] = None, ) -> Dict[str, Any]: """ 키움 ka10001 상세 결과 — 성공/실패 사유를 로그·스크립트용으로 반환. Returns: { "ok": bool, "meta": {"flo_stk": int, "dstr_stk": int} | None, "http_status": int, "return_code": int | None, "return_msg": str, "reason": str, } """ code = str(code or "").strip() empty: Dict[str, Any] = { "ok": False, "meta": None, "http_status": 0, "return_code": None, "return_msg": "", "reason": "", } if not code or not kiwoom_key or not kiwoom_secret: empty["reason"] = "입력값없음(code/키)" return empty token = _get_kiwoom_token_cached(kiwoom_key, kiwoom_secret, is_mock) if not token: empty["reason"] = "키움토큰발급실패" return empty retries = max_retries if retries is None: retries = get_env_int("KIWOOM_KA10001_RETRY_MAX", 4) retry_base = get_env_float("KIWOOM_KA10001_RETRY_SLEEP_SEC", 2.0) domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com" url = f"https://{domain}/api/dostk/stkinfo" headers = { "content-type": "application/json;charset=UTF-8", "appkey": kiwoom_key, "appsecret": kiwoom_secret, "authorization": f"Bearer {token}", "api-id": "ka10001", "cont-yn": "N", "next-key": "", } last: Dict[str, Any] = dict(empty) for attempt in range(max(1, int(retries))): try: resp = requests.post( url, json={"stk_cd": code}, headers=headers, timeout=12, ) http_st = int(resp.status_code) try: data = resp.json() except Exception: data = {} rc_raw = data.get("return_code", 1) try: rc = int(rc_raw) except Exception: rc = 1 msg = str(data.get("return_msg") or "").strip() if http_st == 429 or rc == 5: last = { "ok": False, "meta": None, "http_status": http_st, "return_code": rc, "return_msg": msg, "reason": "rate_limit_1700", } if attempt + 1 < retries: wait = retry_base * (attempt + 1) + random.uniform(0.2, 0.6) logger.warning( "키움 ka10001 레이트리밋 (%s) http=%s rc=%s — %.1fs 후 재시도 %d/%d", code, http_st, rc, wait, attempt + 2, retries, ) time.sleep(wait) continue return last if rc != 0: last = { "ok": False, "meta": None, "http_status": http_st, "return_code": rc, "return_msg": msg, "reason": "api_error", } logger.debug("키움 ka10001 API오류 (%s): rc=%s %s", code, rc, msg) return last flo = int(float(str(data.get("flo_stk", "0") or "0").replace(",", "") or 0)) dstr = int(float(str(data.get("dstr_stk", "0") or "0").replace(",", "") or 0)) if flo <= 0 and dstr <= 0: last = { "ok": False, "meta": None, "http_status": http_st, "return_code": rc, "return_msg": msg, "reason": "주식수없음(flo/dstr=0)", } return last return { "ok": True, "meta": {"flo_stk": flo, "dstr_stk": dstr}, "http_status": http_st, "return_code": rc, "return_msg": msg, "reason": "", } except Exception as e: last = { "ok": False, "meta": None, "http_status": 0, "return_code": None, "return_msg": str(e), "reason": "exception", } if attempt + 1 < retries: time.sleep(retry_base) continue logger.debug("키움 ka10001 예외 (%s): %s", code, e) return last return last def fetch_kiwoom_stock_meta( code: str, kiwoom_key: str, kiwoom_secret: str, is_mock: bool = False, ) -> Optional[Dict[str, int]]: """ 키움 ka10001 주식기본정보 — 상장/유통주식수 (WSManager·stock_share_meta DB). Returns: {"flo_stk": int, "dstr_stk": int} 또는 None """ res = fetch_kiwoom_stock_meta_detail( code, kiwoom_key, kiwoom_secret, is_mock=is_mock, ) if not res.get("ok"): return None return res.get("meta")