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

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

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

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

3278 lines
142 KiB
Python
Raw Permalink Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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)는 스위치(WS_TICK_DB_SAVE_LAG_CUT_ENABLED) 로 제어.
skip_ram = False
_lag_kis = 0.0
try:
from kis_trader.engine.feed_fallback import (
is_feed_read_stale,
packet_lag_seconds,
)
_lag_kis = float(packet_lag_seconds(tick_time_raw or tick_time) or 0.0)
skip_ram = is_feed_read_stale(_lag_kis)
except Exception:
skip_ram = False
# 2026-09-06: DB 저장 컷 스위치 (3벤더 공통 · docs/정합성.md §9)
try:
from database import get_env_bool as _geb, get_env_float as _gef
_kis_db_cut_on = bool(_geb("WS_TICK_DB_SAVE_LAG_CUT_ENABLED", False))
_kis_max_age = float(_gef("LIVE_FEED_FALLBACK_MAX_AGE_SEC", 3.0) or 3.0)
except Exception:
_kis_db_cut_on = False
_kis_max_age = 3.0
_skip_persist_kis = bool(_kis_db_cut_on and _kis_max_age > 0 and _lag_kis > _kis_max_age)
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) 저장
# 2026-09-06: 스위치 OFF(기본) 이면 전부 저장(벤더 통계). ON 이면 lag > LIVE_FEED_FALLBACK 미저장.
if self._tick_recorder is not None and not _skip_persist_kis:
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
last_cleanup_ts = 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)
# ws_candles 구데이터 1시간 주기 정리
if now - last_cleanup_ts >= 3600.0:
last_cleanup_ts = now
try:
keep = get_env_int("SCALP_CANDLE_KEEP_DAYS", 7) or 7
if keep > 0 and self.db is not None and hasattr(self.db, "cleanup_old_ws_candles"):
self.db.cleanup_old_ws_candles(keep_days=keep)
except Exception as e:
logger.debug("ws_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()
# 파일 캐시 경로 설정 (실매매 봇과 백테스트 간 토큰 공유)
import os
base_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
self._cache_file = os.path.join(base_dir, f".kiwoom_token_cache_{'mock' if is_mock else 'real'}.json")
self._load_cache()
def _load_cache(self) -> None:
import os, json
if not os.path.exists(self._cache_file):
return
try:
with open(self._cache_file, "r") as f:
data = json.load(f)
if data.get("app_key") != self._app_key[:10]: # 키가 다르면 무효
return
tk = data.get("token")
exp = data.get("expiry")
if tk and exp:
self._token = tk
self._expiry = _datetime.fromisoformat(exp)
if not self._is_valid():
self._token = None
self._expiry = None
except Exception as e:
logger.warning("키움 토큰 캐시 로드 실패: %s", e)
def _save_cache(self) -> None:
import json
if not self._token or not self._expiry:
return
try:
with open(self._cache_file, "w") as f:
json.dump({
"app_key": self._app_key[:10],
"token": self._token,
"expiry": self._expiry.isoformat()
}, f)
except Exception as e:
logger.warning("키움 토큰 캐시 저장 실패: %s", e)
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"),
)
self._save_cache()
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")