Files
kis_bot/kis_trader/ws/kis_ws.py
2026-07-30 18:05:07 +09:00

2657 lines
114 KiB
Python
Raw 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 폴링을 대체.
보유 종목 코드를 구독해두면 KIS가 체결마다 push → 즉시 가격 캐시 갱신.
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종목 (MAX_STOCKS=3~4 수준에서 문제 없음)
KIS 2026-02-24 경고 준수 (무한 재연결 차단 정책):
- 재연결 횟수 제한: 1시간 내 MAX_RECONNECTS_PER_HOUR 회 초과 시 자동 대기
- 최대 총 재연결 횟수: MAX_RECONNECT_ATTEMPTS 회 초과 시 WebSocket 종료 → REST fallback
- 지수 백오프: 5초 → 10초 → 20초 ... 최대 300초
모의투자(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, Optional, Set
import requests
logger = logging.getLogger("KISWebSocket")
# ------------------------------------------------------------------
# 모듈 수준에서 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_VOLUME = 12 # 하위호환 alias → CNTG_VOL
# 재연결 정책 (KIS 2026-02-24 경고 준수)
MAX_RECONNECT_ATTEMPTS = 10 # 총 재연결 최대 횟수 (STABLE_CONN_RESET_SEC 이상 안정 연결 후 끊기면 초기화)
MAX_RECONNECTS_PER_HOUR = 6 # 1시간 내 재연결 허용 횟수
RECONNECT_BASE_DELAY_SEC = 5.0 # 초기 재연결 대기(초)
RECONNECT_MAX_DELAY_SEC = 300.0 # 최대 재연결 대기(초)
# 이 시간(초) 이상 안정적으로 연결이 유지됐다가 끊기면 카운터를 초기화.
# 예: 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):
self.app_key = app_key
self.app_secret = app_secret
self.is_mock = is_mock
# 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()
# ── 영구 구독 목록 (홀딩 관심종목 등) ──────────────────────
# 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._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
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
# ── 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 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 stop(self, clear_subscriptions: bool = False) -> None:
"""
WebSocket 수신 중단 및 스레드 종료.
Args:
clear_subscriptions: True 시 모든 종목 구독 해제 후 종료
(봇 정상 종료 시 KIS 서버 정리용)
"""
# 옵션: 모든 구독 해제 (KIS 서버 정리용)
if clear_subscriptions and self._connected and self._ws:
with self._sub_lock:
codes = list(self._subscribed)
for code in codes:
self._send_sub_msg(code, subscribe=False)
with self._sub_lock:
self._subscribed.discard(code)
with self._cache_lock:
self._cache.pop(code, None)
logger.info("📡 WebSocket 구독 해제: %s", code)
logger.info("✅ 모든 WebSocket 구독 정리 완료 (%d종목)", len(codes))
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) -> None:
"""
실시간 체결가 구독 등록.
이미 연결 중이면 즉시 구독 메시지 전송, 연결 전이면 연결 성공 시 일괄 등록.
KIS 세션 한도(MAX_SUBSCRIPTIONS=41) 초과 시 등록 거부 후 경고 로그 출력.
"""
code = (code or "").strip()
if not code:
return
with self._sub_lock:
if code in self._subscribed:
return # 이미 구독 중 → 중복 전송 방지
if len(self._subscribed) >= self.MAX_SUBSCRIPTIONS:
logger.warning(
"⚠️ WebSocket 구독 한도 초과(%d/%d) → %s 구독 거부 "
"(KIS 세션 한도 준수: 불필요 종목 구독해제 후 재시도)",
len(self._subscribed), self.MAX_SUBSCRIPTIONS, code,
)
return
self._subscribed.add(code)
if self._connected and self._ws:
self._send_sub_msg(code, subscribe=True)
logger.info("📡 WebSocket 구독 추가: %s (%d/%d)", code, len(self._subscribed), self.MAX_SUBSCRIPTIONS)
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:
self._send_sub_msg(code, subscribe=False)
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: float = 5.0) -> Optional[Dict]:
"""
WebSocket 캐시에서 실시간 가격 + 당일 시고저를 꺼냅니다.
max_age_sec 초 이내 수신된 체결 틱만 유효하게 취급합니다.
반환 형식: {
"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 age > max_age_sec:
# 데이터가 너무 오래됨 → REST로 재확인
return None
return entry.get("data")
@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)
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)
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)
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._ws.close()
except Exception:
pass
def _build_sub_payload(self, code: str, subscribe: bool) -> 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": "H0STCNT0",
"tr_key": code,
}
},
})
def _send_sub_msg(self, code: str, subscribe: bool = True) -> None:
"""WebSocket으로 구독/해제 메시지 전송. 실패 시 조용히 무시."""
if not self._ws:
return
try:
self._ws.send(self._build_sub_payload(code, subscribe))
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] != "H0STCNT0":
return
# 암호화된 데이터는 아직 미지원 (평문만 처리)
if parts[0] == "1":
logger.debug("H0STCNT0 암호화 데이터 수신 (처리 스킵) → REST fallback 권장")
return
# 한 메시지에 여러 건이 포함될 수 있음 (parts[2] = 건수)
# 단순히 parts[3] 전체를 파싱 (단건 기준)
fields = parts[3].split("^")
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"
# 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, # 당일 저가
}
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)
tick_time = fields[self.IDX_TIME].strip() if len(fields) > self.IDX_TIME else ""
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:
use_inc = False
try:
use_inc = bool(self._candle_agg._volume_is_incremental(code))
except Exception:
use_inc = False
agg_vol = cntg_vol if use_inc else acml_vol
self._candle_agg.on_tick(code, price, agg_vol, tick_time)
# TickRecorder: 실매·백테 공통으로 틱당 체결량(CNTG) 저장
if self._tick_recorder is not None:
try:
self._tick_recorder.on_tick(
code, price, cntg_vol, tick_time, source="kis",
)
except Exception as ex:
logger.debug("H0STCNT0→TickRecorder 실패 %s: %s", code, ex)
logger.debug("H0STCNT0 수신: %s%s", code, int(price))
except (ValueError, IndexError) as e:
logger.debug("H0STCNT0 파싱 오류: %s | raw=%s", e, raw[:80])
# ==================================================================
# 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 점유 중에도 세션 양보 — approval 1:1 준수."""
ws = self._ws
if ws is None:
return
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 = []
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
_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 = []
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
_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 = []
self._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
# ── 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:
logger.warning(
"WebSocket approval_key 발급 불가 → %ds 대기 후 재시도",
int(self._reconnect_delay),
)
time.sleep(self._reconnect_delay)
self._reconnect_delay = min(
self._reconnect_delay * 2, self.RECONNECT_MAX_DELAY_SEC
)
continue
# ── 재연결 카운트 및 대기 ─────────────────────────────────
self._reconnect_count += 1
if self._reconnect_count > 1:
self._reconnect_times.append(time.time())
logger.info(
"🔄 WebSocket 재연결 시도 %d/%d (대기 %ds 완료)",
self._reconnect_count,
self.MAX_RECONNECT_ATTEMPTS,
int(self._reconnect_delay),
)
time.sleep(self._reconnect_delay)
self._reconnect_delay = min(
self._reconnect_delay * 2, self.RECONNECT_MAX_DELAY_SEC
)
# ── 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._reconnect_delay = self.RECONNECT_BASE_DELAY_SEC
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:
for i, code in enumerate(codes):
if i > 0:
gap = self._subscribe_gap_sec()
if gap > 0:
time.sleep(gap)
self._send_sub_msg(code, subscribe=True)
if codes:
logger.info(
"📡 WebSocket 구독 일괄 등록: %s (%d종목, gap=%.2f~%.2fs)",
", ".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)),
)
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()
# 핫패스(락 안)에서 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 _volume_is_incremental(self, code: str) -> bool:
return str(code).strip() in self._incremental_vol_codes
def _holding_peak_for(self, code: str) -> Optional[float]:
if not self._peak_provider:
return None
try:
v = self._peak_provider(code)
return float(v) if v and float(v) > 0 else None
except Exception:
return None
# ------------------------------------------------------------------
# 트랙 2: 백그라운드 DB 배치 기록원
# ------------------------------------------------------------------
def _start_db_writer(self) -> None:
"""백그라운드 DB 배치 기록 스레드를 시작합니다."""
self._db_writer_thread = threading.Thread(
target=self._db_writer_loop,
name="CandleDBWriter",
daemon=True, # 봇 종료 시 자동 소멸
)
self._db_writer_thread.start()
logger.info("✅ CandleAggregator DB 기록원 스레드 시작 (배치=%d, 주기=%.1fs)",
self.BATCH_SIZE, self.FLUSH_INTERVAL)
def _db_writer_loop(self) -> None:
"""
배치 기록 루프.
- Queue에서 BATCH_SIZE개 모이면 즉시 배치 INSERT
- BATCH_SIZE 미달이라도 FLUSH_INTERVAL초마다 플러시
- 겸사겸사 봉주기 경과 후 미확정 상태로 방치된 저유동 종목 봉도
같은 주기로 강제확정 점검 (flush_stale_current_candles)
"""
batch: list = []
last_flush = time.time()
last_stale_check = 0.0
stale_check_interval = float(
get_env_int("WS_CANDLE_STALE_CHECK_INTERVAL_SEC", 2)
)
while True:
try:
# FLUSH_INTERVAL 내에서 최대한 많이 모아 배치 구성
timeout = max(0.1, self.FLUSH_INTERVAL - (time.time() - last_flush))
item = self._write_queue.get(timeout=timeout)
if item is None:
# None = 종료 신호
break
batch.append(item)
self._write_queue.task_done()
except queue.Empty:
pass
now = time.time()
should_flush = len(batch) >= self.BATCH_SIZE or \
(batch and (now - last_flush) >= self.FLUSH_INTERVAL)
if should_flush and batch:
self._flush_batch(batch)
batch = []
last_flush = now
if now - last_stale_check >= stale_check_interval:
last_stale_check = now
# 락 밖: env 핫플래그 갱신 (on_tick 경로에서 get_env 금지)
self._refresh_runtime_flags()
try:
self.flush_stale_current_candles()
except Exception as e:
logger.debug("flush_stale_current_candles 실패(무시): %s", e)
# 루프 종료 시 남은 배치 처리
if batch:
self._flush_batch(batch)
logger.info("CandleDBWriter 스레드 종료")
def _refresh_runtime_flags(self) -> None:
"""락 밖에서만 호출. freeze 핫패스 플래그 캐시 갱신."""
try:
self._freeze_on_confirm_cached = bool(
get_env_bool("WS_CANDLE_FREEZE_ON_CONFIRM", True)
)
self._flags_refresh_ts = time.time()
except Exception:
pass
def _ws_candle_freeze_on_confirm(self) -> bool:
"""확정봉 OHLCV 동결 — docs/정합성.md. 기본 true.
⚠ on_tick/_confirm 은 candle_agg._lock 보유 중 호출됨.
여기서 get_env_* 하면 ENV 캐시 만료 시 DB 스냅샷(수 초)이 락을 잡아
전 전략 매수루프가 멈춘다 → 캐시만 반환.
"""
return bool(getattr(self, "_freeze_on_confirm_cached", True))
def _load_confirmed_ohlcv_from_db(
self, code: str, tf: int, candle_times: list,
) -> 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
# IN 절 길이 제한 — 청크
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, holding_peak
FROM ws_candles
WHERE code=%s AND timeframe=%s AND is_confirmed=1
AND candle_time IN ({ph})
""",
(code, int(tf), *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)),
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),
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, holding_peak, updated_at)
VALUES
(%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, 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 = []
for item in batch:
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),
item.get("source", "ws"),
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:
"""
틱 시각(HHMMSS)과 timeframe(분)으로 봉 시작 시각 키를 반환.
날짜는 오늘 날짜를 사용 (장 중에만 동작하므로 안전).
반환 형식: YYYYMMDDHHMM (분 단위, 초 버림)
예) tick_time="091523", timeframe=3 → "202603020915" (09:15봉)
"""
import datetime as _dt
try:
hh = int(tick_time_str[0:2]) if len(tick_time_str) >= 6 else 0
mm = int(tick_time_str[2:4]) if len(tick_time_str) >= 6 else 0
# timeframe 단위로 내림: 3분봉이면 09:17 → 09:15
floored_mm = (mm // timeframe) * timeframe
today = _dt.date.today().strftime("%Y%m%d")
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,
) -> 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)
def _process_tick(self, code: str, price: float, volume: int,
tick_time: str, tf: int, *, market: str = "KR") -> None:
"""
단일 timeframe 에 대한 틱 처리 (lock 내부에서 호출).
[트랙 1] RAM 갱신만 수행, DB 호출 없음 → 블로킹 0ms
[트랙 2] 봉 확정 순간에만 Queue.put_nowait() → 기록원이 비동기 배치 저장
"""
key = (code, tf)
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,
}
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,
}
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, tf = key
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["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 "ws"),
}
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": "ws",
}
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, tf = key
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='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": "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 ""
)
# 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)
with self._lock:
key = (code, tf)
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 = str(row.get("source") or "rest")
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,
}
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 "db")[: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,
}
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
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).
"""
from kis_trader.engine.candle_rollup import rollup_1m_bars_to_tf
tf = int(target_tf)
if tf <= 1:
return 0
with self._lock:
bars_1m = list(self._confirmed.get((code, 1), []))
if not bars_1m:
return 0
rolled = rollup_1m_bars_to_tf(bars_1m, tf)
return self.merge_confirmed_bars(
code, tf, rolled, log_tag=f"rollup_1m→{tf}M",
)
# ------------------------------------------------------------------
# [트랙 1] RAM 버퍼 조회 — 매수/매도 루프에서 직접 호출 (DB 조회 없음)
# ------------------------------------------------------------------
def get_latest_confirmed(self, code: str, tf: int) -> Optional[dict]:
"""
가장 최근 확정된 봉(완성된 마지막 봉)을 반환.
None이면 아직 봉이 확정되지 않음 (장 초반 등).
"""
with self._lock:
buf = self._confirmed.get((code, tf))
return buf[-1] if buf else None
def get_prev_confirmed(self, code: str, tf: int) -> Optional[dict]:
"""직전 확정봉 (최신에서 2번째). 패턴 확인용 (현재봉 - 1)."""
with self._lock:
buf = self._confirmed.get((code, tf))
return buf[-2] if buf and len(buf) >= 2 else None
def get_candles(self, code: str, tf: int, n: int = 10) -> list:
"""최근 n개 확정 봉 리스트 반환 (오래된→최신 순)."""
with self._lock:
buf = self._confirmed.get((code, tf), [])
return list(buf[-n:])
def get_confirmed_count(self, code: str, tf: int) -> int:
"""확정된 봉 수 (RSI 안정화 여부 확인용)."""
with self._lock:
return len(self._confirmed.get((code, tf), []))
def get_current_candle(self, code: str, tf: int) -> Optional[dict]:
"""
현재 진행 중인 봉(미확정, is_confirmed=0) 반환.
RSI는 포함되지 않음 (확정 봉 기준으로만 계산).
"""
with self._lock:
return dict(self._current.get((code, tf), {})) or None
def get_rsi(self, code: str, tf: int, period: int = 3) -> Optional[float]:
"""
최신 확정 봉의 RSI(period) 값 반환.
period: 2, 3, 5 중 하나 (스캘핑 단타용 초단기 RSI)
"""
candle = self.get_latest_confirmed(code, tf)
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 tf in list(self.timeframes):
key = (code, tf)
self._confirmed.pop(key, None)
self._closes.pop(key, None)
self._current.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_is_rate_limit(http_st: int, rc: Any, msg: str) -> bool:
if int(http_st or 0) == 429:
return True
try:
if int(rc) == 5:
return True
except Exception:
pass
m = str(msg or "")
return ("1700" in m) or ("허용된 요청" in m) or ("유량" in m)
class KiwoomTokenManager:
"""
kiwoom_rest_api/auth/token.py TokenManager 와 동일한 로직 — DB 키 직접 주입 버전.
- get_token(): 유효한 토큰이면 바로 반환, 만료 시에만 재발급 (au10001 방지)
- expires_dt 를 정확히 파싱 (하드코딩 23h 아님)
- 30초 버퍼: 경계 케이스 방지 (TokenManager 와 동일)
- 스레드 안전: 인스턴스당 Lock 보유
"""
def __init__(self, app_key: str, app_secret: str, is_mock: bool = False):
self._app_key = app_key
self._app_secret = app_secret
self._is_mock = is_mock
self._domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com"
self._mode_str = "모의" if is_mock else "실전"
self._token: Optional[str] = None
self._expiry: Optional[_datetime] = None
self._lock = _threading.Lock()
def _is_valid(self) -> bool:
"""만료 30초 전까지 유효 (TokenManager._is_access_token_valid 와 동일)"""
if not self._token or not self._expiry:
return False
return _datetime.now() < self._expiry - _timedelta(seconds=30)
def _request_new_token(self) -> None:
"""토큰 신규 발급 — 만료 시에만 호출됨"""
resp = requests.post(
f"https://{self._domain}/oauth2/token",
json={
"grant_type": "client_credentials",
"appkey": self._app_key,
"secretkey": self._app_secret,
},
timeout=10,
)
data = resp.json()
token = (data.get("token") or data.get("access_token") or "").strip()
if not token:
raise RuntimeError(f"키움 토큰 발급 실패 [{self._mode_str}]: {data}")
# expires_dt 정확히 파싱 (TokenManager._update_token_info 와 동일 로직)
exp_s = data.get("expires_dt", "")
try:
self._expiry = _datetime.strptime(str(exp_s), "%Y%m%d%H%M%S")
except Exception:
# expires_in 도 없으면 24h 기본 (보수적 fallback)
exp_in = data.get("expires_in", 86400)
self._expiry = _datetime.now() + _timedelta(seconds=int(exp_in))
self._token = token
logger.info(
"✅ 키움 토큰 발급 완료 [%s] (앞8자: %s…, 만료: %s)",
self._mode_str, token[:8], self._expiry.strftime("%Y-%m-%d %H:%M:%S"),
)
def invalidate(self, reason: str = "") -> None:
"""서버가 토큰을 거부했을 때 로컬 캐시를 비운다 (만료 전 재사용 금지).
WS LOGIN rc=805004 등 — expires_dt 전이라도 서버 측 무효면
같은 토큰을 재전송하면 재연결만 소모하고 영구 비활성으로 간다.
"""
with self._lock:
self._token = None
self._expiry = None
if reason:
logger.warning(
"♻️ 키움 토큰 캐시 무효화 [%s]: %s",
self._mode_str, reason,
)
def get_token(self) -> Optional[str]:
"""
유효한 토큰 반환. 만료 시에만 재발급.
kiwoom_rest_api TokenManager.get_token() 호환 인터페이스.
"""
with self._lock:
if not self._is_valid():
try:
self._request_new_token()
except Exception as e:
logger.warning("⚠️ 키움 토큰 발급 예외 [%s]: %s", self._mode_str, e)
return None
return self._token
def _kiwoom_token_cache_key(kiwoom_key: str, is_mock: bool) -> str:
domain = "mockapi.kiwoom.com" if is_mock else "api.kiwoom.com"
return f"{domain}:{str(kiwoom_key or '')[:8]}"
def _get_kiwoom_token_cached(
kiwoom_key: str,
kiwoom_secret: str,
is_mock: bool,
) -> Optional[str]:
"""
KiwoomTokenManager 싱글톤 풀에서 인스턴스를 꺼내 토큰 반환.
동일 키 조합은 같은 인스턴스를 재사용 → au10001 rate limit 방지.
"""
cache_key = _kiwoom_token_cache_key(kiwoom_key, is_mock)
with _kiwoom_token_lock:
if cache_key not in _kiwoom_managers:
_kiwoom_managers[cache_key] = KiwoomTokenManager(
kiwoom_key, kiwoom_secret, is_mock=is_mock
)
mgr = _kiwoom_managers[cache_key]
return mgr.get_token()
def _invalidate_kiwoom_token_cached(
kiwoom_key: str,
kiwoom_secret: str,
is_mock: bool,
reason: str = "",
) -> None:
"""LOGIN/REST 토큰 거부 시 싱글톤 캐시 무효화 — 다음 get_token 이 재발급."""
cache_key = _kiwoom_token_cache_key(kiwoom_key, is_mock)
with _kiwoom_token_lock:
mgr = _kiwoom_managers.get(cache_key)
if mgr is not None:
mgr.invalidate(reason or "token rejected")
def _get_kiwoom_creds(db) -> tuple:
"""
DB env_config 최신 행에서 키움 앱키/시크릿 반환.
KIS_MOCK 설정에 따라 MOCK / REAL 키를 자동 선택.
Returns:
(app_key, app_secret, is_mock) — 키 없으면 (None, None, False)
"""
try:
row = db.conn.execute(
"SELECT * FROM env_config ORDER BY id DESC LIMIT 1"
).fetchone()
if not row:
return None, None, False
r = dict(row)
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):
sem.acquire()
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_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")