- 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>
3278 lines
142 KiB
Python
3278 lines
142 KiB
Python
"""
|
||
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")
|
||
|