브랜치 분리 방식: A / B / C

A 선택 시 커밋 메시지: 위 초안 OK / 수정 / 직접 작성
작업 시점: 지금 / 운영 데이터 1~2일 쌓고 / 주말
This commit is contained in:
2026-05-05 21:04:17 +09:00
parent c2b2b711e0
commit f61c471aac
58 changed files with 803502 additions and 1430 deletions

View File

@@ -0,0 +1,4 @@
"""kis_trader.network — 통합 WebSocket/Event Bus 허브."""
from .ws_manager import WSManager
__all__ = ["WSManager"]

View File

@@ -0,0 +1,396 @@
"""
kis_trader/network/condition_manager.py — KIS 조건검색 기반 동적 유니버스
===========================================================================
팩트 체크 먼저:
* KIS 는 ``H0UPANC0`` 웹소켓으로 "조건검색 실시간" 을 주지 않는다.
H0UPANC0 는 **업종별 예상체결** TR 이다. 인터넷 블로그/LLM 답변에
자주 보이는 "조건검색 웹소켓" 은 대부분 키움 OpenAPI+ 쪽 이야기.
* KIS 공식 경로는 REST 두 개뿐:
/quotations/psearch-title → 서버 저장 조건식 목록 (HHKST03900300)
/quotations/psearch-result → 특정 조건식 현재 결과 (HHKST03900400)
* 따라서 REST 폴링이 유일한 방법. 기본 폴링 주기는 ``CONDITION_POLL_INTERVAL_SEC``
(기본 10초). 10초면 종목당 하루 ~2,340 호출로 429 안전 여유 충분.
v2 변경점 (다중 조건식 지원):
* 전략마다 다른 조건식을 쓸 수 있도록 ``configs`` 인자 추가:
configs = [
{"strategy_id": "SCALP", "name": "체결강도급등", "seq": "0"},
{"strategy_id": "SHORT", "name": "꼬리달린봉", "seq": ""}, # seq 는 name 으로 자동 해결
{"strategy_id": "BREAKOUT", "name": "우상향돌파", "seq": ""},
]
* 같은 조건식을 여러 전략이 공유해도 OK (seq 가 같으면 REST 1번만 호출)
* 전략별 get_universe_for / get_candidates_for 제공
* **변동(ENTER/EXIT) 감지 tick 마다** ``target_candidates_history`` 에
초단위 ``event_time`` (YYYY-MM-DD HH:MM:SS) 으로 풀 스냅샷을 INSERT.
(strategy_id 컬럼은 ``TradeDBExt.insert_condition_universe_snapshot``
이 자동 마이그레이션.) 백테스트는 ``TradeDBExt.get_universe_by_candle_time()``
으로 "그 1분봉 시점에 봇이 보던 유니버스" 를 재현.
사용 (권장 — multi):
cm = ConditionSearchManager(
client=kis_client,
user_id="HTSID",
configs=[
{"strategy_id": "SCALP", "name": "체결강도급등"},
{"strategy_id": "BREAKOUT", "seq": "0"},
],
db=db,
)
cm.start()
codes: set = cm.get_universe_for("SCALP")
사용 (legacy — single, 기존 호출 호환):
cm = ConditionSearchManager(
client=kis_client, user_id="HTSID", condition_name="우상향돌파",
)
"""
from __future__ import annotations
import random
import threading
import time
from datetime import datetime as dt
from typing import Callable, Dict, List, Optional, Set
from ..utils.env import get_env_bool, get_env_int
from ..utils.logger import get_logger
logger = get_logger("kis_trader.cond")
class ConditionSearchManager:
"""KIS 조건검색 폴링 매니저. 여러 조건식을 동시에 병렬 관리."""
def __init__(
self,
*,
client,
user_id: str,
configs: Optional[List[Dict]] = None,
condition_name: Optional[str] = None,
condition_seq: Optional[str] = None,
on_change: Optional[Callable[[str, Set[str], Set[str], Set[str]], None]] = None,
poll_interval_sec: Optional[float] = None,
db=None,
):
self.client = client
self.user_id = (user_id or "").strip()
self.on_change = on_change
self.db = db # 히스토리 저장용 (선택)
self.poll_interval = float(
poll_interval_sec
if poll_interval_sec is not None
else get_env_int("CONDITION_POLL_INTERVAL_SEC", 10)
)
self.history_enabled = get_env_bool("CONDITION_HISTORY_SAVE", True)
# EXIT grace period: 한 번 빠진 종목을 N초간 universe 에 keep.
# 단발성 EXIT/RE-ENTER 회전을 흡수해 WS 구독 해제 → 캐시·갭보정 리셋
# → resubscribe 후 데이터 부족으로 매수 시그널 못내는 사이클을 차단.
# 기본 0 = 비활성 (기존 동작과 동일). 권장 60.
self._exit_grace_sec = float(get_env_int("CONDITION_EXIT_GRACE_SEC", 0))
# strategy_id → {code: first_missing_at_epoch}
self._pending_exit: Dict[str, Dict[str, float]] = {}
# 설정 정규화: legacy(single) → multi 형식으로 흡수
self._configs: List[Dict] = []
if configs:
for c in configs:
sid = str(c.get("strategy_id") or "").strip().upper()
nm = (c.get("name") or "").strip() or None
sq = (c.get("seq") or "").strip() or None
if not sid or (not nm and not sq):
continue
self._configs.append({"strategy_id": sid, "name": nm, "seq": sq})
elif condition_name or condition_seq:
self._configs.append({
"strategy_id": "DEFAULT",
"name": (condition_name or "").strip() or None,
"seq": (condition_seq or "").strip() or None,
})
self._thread: Optional[threading.Thread] = None
self._running = False
# 상태
self._current: Dict[str, Set[str]] = {} # strategy_id → code set
self._name_map: Dict[str, str] = {} # code → name (전역)
# 초기 tick 에서 "빈 set → 첫 결과" 를 변동으로 간주해 1회는 저장
self._initialized: Set[str] = set()
self._lock = threading.Lock()
# ------------------------------------------------------------------
# Public API
# ------------------------------------------------------------------
def start(self) -> bool:
"""조건식 seq 를 해결하고 폴링 쓰레드 기동. 유효 조건식 0개면 False."""
if not self.user_id:
logger.warning("조건검색 user_id 누락 → 비활성")
return False
if not self._configs:
logger.info("조건검색 configs 비어 있음 → 매니저 비활성")
return False
# seq 해결 (이름 → seq). 1번만 전체 목록 호출해서 캐시.
name_to_seq = self._fetch_seq_map()
valid = []
for cfg in self._configs:
if not cfg.get("seq") and cfg.get("name"):
sq = name_to_seq.get(cfg["name"])
if sq:
cfg["seq"] = sq
if cfg.get("seq"):
valid.append(cfg)
logger.info(
"🔗 조건식 매핑: strategy=%s seq=%s name=%s",
cfg["strategy_id"], cfg["seq"], cfg.get("name") or "?",
)
else:
logger.warning(
"⚠️ 조건식 seq 해결 실패 (strategy=%s name=%s) → 이 전략은 폴백",
cfg["strategy_id"], cfg.get("name"),
)
self._configs = valid
if not self._configs:
logger.warning("유효 조건식 0개 → 조건검색 매니저 비활성")
return False
self._running = True
self._thread = threading.Thread(
target=self._loop, daemon=True, name="CondSearch"
)
self._thread.start()
logger.info(
"✅ 조건검색 폴링 시작 (%d개, interval=%ds, exit_grace=%ds, history=%s)",
len(self._configs), int(self.poll_interval), int(self._exit_grace_sec),
"ON" if (self.history_enabled and self.db is not None) else "OFF",
)
return True
def stop(self) -> None:
self._running = False
# ── 전략별 조회 ────────────────────────────────────────────
def get_universe_for(self, strategy_id: str) -> Set[str]:
sid = (strategy_id or "").upper()
with self._lock:
return set(self._current.get(sid, set()))
def get_candidates_for(self, strategy_id: str) -> List[Dict]:
"""BaseStrategy._load_candidates 와 호환되는 dict 리스트 반환."""
sid = (strategy_id or "").upper()
with self._lock:
codes = list(self._current.get(sid, set()))
nm = dict(self._name_map)
# 전략별 기본 필터 플래그는 True 로 열어둔다 (전략 쪽 _candidate_filter 가 판단)
out = []
for c in codes:
out.append({
"code": c,
"name": nm.get(c, c),
"scalp_on": True,
"tail_on": True,
"score": 0.0,
"price": 0.0,
})
return out
# ── 하위 호환 (Breakout 단일 매니저용) ───────────────────
def get_universe(self) -> Set[str]:
"""모든 전략 유니버스의 합집합 (heartbeat/총량 로그용)."""
with self._lock:
out: Set[str] = set()
for s in self._current.values():
out |= s
return out
def get_candidates(self) -> List[Dict]:
"""기본 호출 시 첫 번째 설정된 전략 유니버스 반환 (하위호환)."""
if not self._configs:
return []
return self.get_candidates_for(self._configs[0]["strategy_id"])
# ------------------------------------------------------------------
# 내부
# ------------------------------------------------------------------
@staticmethod
def _row_name(row: Dict) -> str:
"""
KIS psearch-title 응답의 '조건식 이름' 추출.
실제 응답 키는 ``condition_nm`` (예: {"seq":"0","condition_nm":"돌파_초반강세",...}).
과거/문서상 변형 키도 모두 허용해 안전하게 폴백.
"""
for k in ("condition_nm", "condition_name", "cond_nm", "user_cnd_nm"):
v = row.get(k)
if v is not None and str(v).strip():
return str(v).strip()
return ""
def _fetch_seq_map(self) -> Dict[str, str]:
"""서버 저장 조건식 목록 1회 호출 → name→seq 맵."""
try:
lst = self.client.get_condition_list(self.user_id) or []
except Exception as e:
logger.error("조건식 목록 조회 예외: %s", e)
return {}
if not lst:
logger.warning("조건식 목록이 비어있음 (user_id=%s)", self.user_id)
return {}
logger.info(
"저장된 조건식 %d개: %s",
len(lst),
", ".join(f"{x.get('seq')}:{self._row_name(x) or '?'}" for x in lst),
)
return {
self._row_name(row): str(row.get("seq") or "").strip()
for row in lst
if self._row_name(row)
}
def _loop(self) -> None:
while self._running:
try:
self._tick_all()
except Exception as e:
logger.error("조건검색 루프 예외: %s", e)
# 서버 부하 방지 지터 — 주기 대비 10% 수준
jitter = min(1.0, self.poll_interval * 0.1)
sleep_sec = self.poll_interval + random.uniform(0, jitter)
# 중단 감지 해상도 0.5s (10초 주기에 1초 해상도는 과함)
deadline = time.time() + sleep_sec
while time.time() < deadline:
if not self._running:
return
time.sleep(0.5)
def _tick_all(self) -> None:
# 같은 seq 를 공유하는 전략이 있으면 REST 1회만 호출 (캐시)
seq_cache: Dict[str, List[Dict]] = {}
for cfg in self._configs:
sid = cfg["strategy_id"]
seq = cfg["seq"]
if seq in seq_cache:
rows = seq_cache[seq]
else:
try:
rows = self.client.get_condition_result(self.user_id, seq) or []
except Exception as e:
logger.debug("조건검색 결과 조회 실패 (%s/%s): %s", sid, seq, e)
rows = []
seq_cache[seq] = rows
self._apply_result(sid, rows)
def _apply_result(self, strategy_id: str, rows: List[Dict]) -> None:
"""
조건검색 raw 결과를 universe 에 반영.
EXIT grace 정책 (``CONDITION_EXIT_GRACE_SEC`` > 0 일 때):
- raw 결과에서 빠진 종목을 즉시 EXIT 처리하지 않고 ``_pending_exit`` 에
``first_missing_at`` 시각과 함께 등록.
- grace 초가 지나야 진짜 EXIT (universe 에서 제거) → WS 구독 해제 → 갭보정 리셋.
- grace 중 다시 raw 에 등장하면 pending 에서 빼고 universe 에 그대로 keep
(구독·캐시·봉 데이터 보존).
효과: 단발성 EXIT/RE-ENTER 폭주 흡수. 시장 노이즈로 1~2 tick 빠지는 케이스를
걸러 매수 시그널 평가용 RAM 데이터 (회복률·낙폭 등) 가 휘발 안 됨.
"""
raw_set = {r["code"] for r in rows if r.get("code")}
new_names = {
r["code"]: r.get("name", r["code"]) for r in rows if r.get("code")
}
now = time.time()
grace = self._exit_grace_sec
with self._lock:
prev = self._current.get(strategy_id, set())
first_tick = strategy_id not in self._initialized
pending = self._pending_exit.setdefault(strategy_id, {})
# ── grace 적용 ────────────────────────────────────────
kept_in_grace: Set[str] = set()
if grace > 0:
# 1) raw 에 다시 등장 → grace 해제 (회생)
for c in raw_set & set(pending.keys()):
pending.pop(c, None)
# 2) raw 에서 빠진 prev 종목 → pending 등록 (처음 사라진 시각)
for c in prev - raw_set:
if c not in pending:
pending[c] = now
# 3) grace 미경과 종목 → universe 에 keep / 경과 → pending 제거
expired: Set[str] = set()
for c, first_at in list(pending.items()):
if now - first_at >= grace:
expired.add(c)
for c in expired:
pending.pop(c, None)
kept_in_grace = set(pending.keys())
# grace=0 (기존 동작) 이면 kept_in_grace = ∅, pending 안 쓰임.
# 봇이 실제로 보는 effective universe (raw + grace keep)
new_set = raw_set | kept_in_grace
# ── 변동 계산 (effective 기준) ───────────────────────
enters = new_set - prev
exits = prev - new_set
self._current[strategy_id] = new_set
self._initialized.add(strategy_id)
for c, n in new_names.items():
self._name_map[c] = n
changed = bool(enters or exits)
n_grace = len(kept_in_grace)
if changed:
grace_tag = f", grace={n_grace}" if n_grace else ""
logger.info(
"🔄 [%s] +%d / -%d (현재 %d종목%s)",
strategy_id, len(enters), len(exits), len(new_set), grace_tag,
)
if enters:
preview = ", ".join(
f"{c}({new_names.get(c, c)})" for c in list(sorted(enters))[:5]
)
logger.info(" ENTER: %s%s",
preview, "" if len(enters) > 5 else "")
if exits:
preview = ", ".join(sorted(exits)[:5])
logger.info(" EXIT : %s%s",
preview, "" if len(exits) > 5 else "")
if self.on_change and changed:
try:
self.on_change(strategy_id, new_set, enters, exits)
except Exception as e:
logger.warning("on_change 콜백 예외: %s", e)
# 백테스트 재현성 보장: 첫 tick 또는 변동 발생 tick 마다 풀 스냅샷 저장
# (effective universe 기준 — 봇이 실제로 보던 universe 가 그대로 기록됨)
# 변동 없는 tick 은 공간 절약 위해 skip
if changed or first_tick:
self._save_snapshot(strategy_id, new_set, new_names)
def _save_snapshot(
self,
strategy_id: str,
codes: Set[str],
names: Dict[str, str],
) -> None:
"""변동이 감지된 tick 의 풀 유니버스 스냅샷 저장."""
if not (self.history_enabled and self.db is not None):
return
# 유니버스가 비어있어도 "비었다" 는 사실을 기록해야 백테스트에서 재현 가능.
# 단, 한 번도 결과를 못 받은 상태(첫 호출 실패 등)는 저장 X.
event_time = dt.now().strftime("%Y-%m-%d %H:%M:%S")
items = [{"code": c, "name": names.get(c, c)} for c in sorted(codes)]
try:
n = self.db.insert_condition_universe_snapshot(
strategy_id=strategy_id,
event_time=event_time,
items=items,
)
logger.debug(
"📼 [history] %s @%s %d종목 저장",
strategy_id, event_time, n,
)
except Exception as e:
logger.debug("history 저장 예외: %s", e)

View File

@@ -0,0 +1,330 @@
"""
kis_trader/network/market_guard.py — 시장 급락 서킷브레이커
==============================================================
KOSPI/KOSDAQ 종합지수를 주기적으로 폴링해, 다음 조건 충족 시 PANIC 모드 진입:
1) 5분 누적 -N% 하락 (기본 -2%)
2) 일중 누적 -N% 하락 (전일 종가 대비, 기본 -3%)
PANIC 모드 시:
- ``BaseStrategy._scan_and_buy()`` 가 신규 매수를 즉시 차단
- 기존 보유 종목 매도는 정상 동작 (포지션 정리·손실 확대 방지)
해제:
- 5분 누적 +N% 반등 시 자동 해제 (기본 +1%)
- 또는 운영자가 DB ``MARKET_GUARD_ENABLED=false`` 로 수동 해제
설계 철학
---------
- 거래소 공식 서킷브레이커 (KOSPI -8%) 보다 훨씬 빨리 반응 → 봇 보호 우선.
- 매수만 차단, 매도는 평소처럼 진행 → 봇이 살아있어야 손절·익절 가능.
- 백테스트 X (지수 데이터 누적 안 되어 있음) → 실거래로 임계치 튜닝.
- ``MARKET_GUARD_ENABLED=false`` 가 기본값. 운영 1~2주 모니터링 후 활성화.
환경변수 (DB env_config, 매 tick 재조회 → 운영 중 즉시 반영)
------------------------------------------------------------
MARKET_GUARD_ENABLED (true/false, 기본 false)
MARKET_GUARD_5MIN_DROP_PCT (실수, 기본 2.0)
MARKET_GUARD_DAILY_DROP_PCT (실수, 기본 3.0)
MARKET_GUARD_RECOVERY_PCT (실수, 기본 1.0)
MARKET_GUARD_INDEX_CODE ("0001"|"1001"|"both", 기본 "both")
MARKET_GUARD_POLL_SEC (정수, 기본 30)
"""
from __future__ import annotations
import threading
import time
from collections import deque
from datetime import datetime as dt
from typing import Deque, Dict, List, Optional, Tuple
from ..utils.env import get_env_bool, get_env_float, get_env_from_db, get_env_int
from ..utils.logger import get_logger
logger = get_logger("kis_trader.market_guard")
class MarketGuard:
"""KOSPI/KOSDAQ 지수 폴링 + PANIC 판정 매니저. 독립 백그라운드 쓰레드."""
# 지수 코드 → 사람용 이름
_INDEX_NAMES = {
"0001": "KOSPI",
"1001": "KOSDAQ",
"2001": "KOSPI200",
}
def __init__(self, *, client, db=None):
self.client = client
self.db = db
self._thread: Optional[threading.Thread] = None
self._running = False
self._lock = threading.Lock()
# 지수별 최근 가격 deque [(epoch_sec, value), ...] — 5분 비교용
# 보수적으로 10분치 (deque maxlen=120 ≈ poll 5초 × 120 = 10분)
self._history: Dict[str, Deque[Tuple[float, float]]] = {}
# 지수별 전일 종가 (일중 누적 판정용)
self._prev_close: Dict[str, float] = {}
# PANIC 상태 (전역 — 어떤 지수가 트리거했든 매수 전체 차단)
self._panic = False
self._panic_reason: str = ""
self._panic_since: float = 0.0
self._panic_index: str = ""
self._reload_config()
# ------------------------------------------------------------------
def _reload_config(self) -> None:
"""매 tick DB 재조회 (운영 중 임계치 튜닝 즉시 반영)."""
self.enabled = get_env_bool("MARKET_GUARD_ENABLED", False)
self.drop_5min = get_env_float("MARKET_GUARD_5MIN_DROP_PCT", 2.0)
self.drop_daily = get_env_float("MARKET_GUARD_DAILY_DROP_PCT", 3.0)
self.recovery_pct = get_env_float("MARKET_GUARD_RECOVERY_PCT", 1.0)
self.poll_sec = max(5, get_env_int("MARKET_GUARD_POLL_SEC", 30))
# 감시 지수: "0001", "1001", "both", 또는 콤마구분 ("0001,1001")
idx_raw = (get_env_from_db("MARKET_GUARD_INDEX_CODE", "both") or "both") \
.strip().lower()
if idx_raw == "both":
self.index_codes = ["0001", "1001"]
else:
codes = [c.strip() for c in idx_raw.split(",") if c.strip()]
# 알 수 없는 코드 들어오면 안전하게 KOSPI 만 감시
self.index_codes = codes if codes else ["0001"]
# ------------------------------------------------------------------
def start(self) -> bool:
"""백그라운드 쓰레드 기동. 비활성 상태로 시작해도 thread 자체는 살아있음
→ DB 토글로 즉시 활성화 가능."""
if self._thread and self._thread.is_alive():
return True
self._running = True
self._thread = threading.Thread(
target=self._loop, daemon=True, name="MarketGuard"
)
self._thread.start()
logger.info(
"✅ MarketGuard 시작 (enabled=%s, indexes=%s, "
"5min=-%.1f%%, daily=-%.1f%%, recovery=+%.1f%%, poll=%ds)",
self.enabled, self.index_codes, self.drop_5min, self.drop_daily,
self.recovery_pct, self.poll_sec,
)
return True
def stop(self) -> None:
self._running = False
# ------------------------------------------------------------------
# 외부 API — 전략에서 호출
# ------------------------------------------------------------------
def is_panic(self) -> bool:
"""매수 차단 여부. BaseStrategy._scan_and_buy() 가 매 후보마다 호출."""
with self._lock:
return self.enabled and self._panic
def panic_reason(self) -> str:
"""현재 PANIC 사유 (로그용)."""
with self._lock:
return self._panic_reason
def status(self) -> Dict:
"""heartbeat 로그 등에서 노출용 — 현재 상태 요약."""
with self._lock:
now = time.time()
indexes_status = []
for code in self.index_codes:
name = self._INDEX_NAMES.get(code, code)
hist = self._history.get(code)
prdy = self._prev_close.get(code, 0)
if not hist or prdy <= 0:
indexes_status.append({"name": name, "ready": False})
continue
cur = hist[-1][1]
daily_pct = (cur - prdy) / prdy * 100.0
five = self._lookup_price_at(hist, now - 300)
five_pct = ((cur - five) / five * 100.0) if five else 0.0
indexes_status.append({
"name": name, "ready": True,
"daily_pct": round(daily_pct, 2),
"five_min_pct": round(five_pct, 2),
})
return {
"enabled": self.enabled,
"panic": self._panic,
"reason": self._panic_reason,
"since_min": round((now - self._panic_since) / 60.0, 1) if self._panic else 0.0,
"indexes": indexes_status,
}
# ------------------------------------------------------------------
# 메인 루프
# ------------------------------------------------------------------
def _loop(self) -> None:
last_log_ts = 0.0
while self._running:
try:
self._reload_config()
# 비활성 시 천천히 대기 (DB 폴링 부하 ↓, 토글 즉시 반응)
if not self.enabled:
time.sleep(max(10, self.poll_sec))
continue
# 장중에만 감시 (장외에는 지수 데이터 정적 → 폴링 무의미)
if not self._is_market_hours():
time.sleep(60)
continue
self._tick()
# 60초에 1번 상태 로그 (정상 동작 확인용)
now = time.time()
if now - last_log_ts >= 60:
last_log_ts = now
self._log_status()
except Exception as e:
logger.error("MarketGuard 루프 예외: %s", e)
time.sleep(self.poll_sec)
@staticmethod
def _is_market_hours() -> bool:
"""장 시간 체크 (BaseStrategy.check_market_status 와 동일 정책)."""
if get_env_bool("FORCE_MARKET_OPEN", False):
return True
now = dt.now()
h, m = now.hour, now.minute
return (9 <= h < 15) or (h == 15 and m <= 30)
# ------------------------------------------------------------------
def _tick(self) -> None:
"""모든 감시 지수 1회 폴링 + PANIC 판정."""
now = time.time()
for code in self.index_codes:
try:
out = self.client.inquire_index_price(code)
except Exception as e:
logger.debug("지수 조회 실패 (%s): %s", code, e)
continue
if not out:
continue
try:
cur = float(str(out.get("bstp_nmix_prpr", 0)).replace(",", ""))
prdy_close = float(
str(out.get("bstp_nmix_prdy_clpr", 0)).replace(",", "")
)
except (ValueError, AttributeError):
continue
if cur <= 0 or prdy_close <= 0:
continue
with self._lock:
# 전일 종가 저장 (일중 누적 판정용 — 매번 갱신 OK, 같은 값)
self._prev_close[code] = prdy_close
# 가격 deque 업데이트 (10분 이상 된 데이터 자동 삭제)
hist = self._history.setdefault(code, deque(maxlen=120))
hist.append((now, cur))
cutoff = now - 600
while hist and hist[0][0] < cutoff:
hist.popleft()
# ── PANIC 판정 ────────────────────────────────
self._evaluate_panic(code, cur, prdy_close, now)
def _evaluate_panic(
self, code: str, cur: float, prdy_close: float, now: float,
) -> None:
"""단일 지수 기준 PANIC 진입/해제 판정. _lock 보유 상태 가정."""
name = self._INDEX_NAMES.get(code, code)
# 일중 누적 등락률 (전일 종가 대비)
daily_chg = (cur - prdy_close) / prdy_close * 100.0
# 5분 누적 등락률 (5분 전 가격 vs 현재가)
hist = self._history.get(code, deque())
five_min_ago_price = self._lookup_price_at(hist, now - 300)
if five_min_ago_price:
five_min_chg = (cur - five_min_ago_price) / five_min_ago_price * 100.0
else:
five_min_chg = 0.0 # 5분치 데이터 아직 부족 (시작 직후)
if not self._panic:
# PANIC 진입 조건
reasons: List[str] = []
if daily_chg <= -abs(self.drop_daily):
reasons.append(f"일중 {daily_chg:+.2f}% (한계 -{self.drop_daily:.1f}%)")
if five_min_chg <= -abs(self.drop_5min):
reasons.append(f"5분 {five_min_chg:+.2f}% (한계 -{self.drop_5min:.1f}%)")
if reasons:
self._panic = True
self._panic_index = code
self._panic_reason = f"{name} " + " / ".join(reasons)
self._panic_since = now
logger.warning(
"🚨 [MarketGuard] PANIC 진입: %s — 모든 전략 신규 매수 차단",
self._panic_reason,
)
else:
# PANIC 해제 조건: 트리거된 지수의 5분 누적 +N% 반등
# (다른 지수 반등은 무시 — 같은 지수가 회복해야 진짜 회복)
if code != self._panic_index:
return
if five_min_chg >= abs(self.recovery_pct):
duration_min = (now - self._panic_since) / 60.0
logger.warning(
"✅ [MarketGuard] PANIC 해제: %s 5분 %+.2f%% 반등 "
"(지속 %.1f분, 매수 재개)",
name, five_min_chg, duration_min,
)
self._panic = False
self._panic_reason = ""
self._panic_index = ""
self._panic_since = 0.0
@staticmethod
def _lookup_price_at(
hist: Deque[Tuple[float, float]], target_ts: float,
) -> Optional[float]:
"""deque 에서 target_ts 이전(≤) 시각 중 가장 최근 가격 반환.
없으면 None (= 데이터 부족, 5분 비교 스킵)."""
best = None
for ts, val in hist:
if ts <= target_ts:
best = val
else:
break
return best
def _log_status(self) -> None:
"""60초에 한 번씩 평상시/PANIC 상태 로그."""
with self._lock:
if self._panic:
duration_min = (time.time() - self._panic_since) / 60.0
logger.warning(
"⛔ [MarketGuard] PANIC 모드 진행 중 — %s (지속 %.1f분)",
self._panic_reason, duration_min,
)
return
# 평상시 한 줄 요약 (감시 지수별 일중/5분 등락률)
parts = []
now = time.time()
for code in self.index_codes:
name = self._INDEX_NAMES.get(code, code)
hist = self._history.get(code)
prdy = self._prev_close.get(code, 0)
if not hist or prdy <= 0:
parts.append(f"{name}=대기")
continue
cur = hist[-1][1]
daily = (cur - prdy) / prdy * 100.0
five = self._lookup_price_at(hist, now - 300)
five_chg = ((cur - five) / five * 100.0) if five else 0.0
parts.append(f"{name} 일중{daily:+.2f}% 5분{five_chg:+.2f}%")
if parts:
logger.info("📊 [MarketGuard] 정상 — %s", " | ".join(parts))

View File

@@ -0,0 +1,313 @@
"""
kis_trader/network/ranking_manager.py — 거래량/체결강도 순위 기반 유니버스
============================================================================
전략별 유니버스를 KIS 랭킹 REST(volume-rank FHPST01710000) 로 갱신하는 매니저.
왜 조건검색 대신 랭킹인가:
* 조건검색은 "구조적 필터" (20일 신고가·이평정배열 등 돌파용) 에 강하지만
"지금 반등 중" 같은 동적 판단엔 약하다.
* 실제 매매 시점 판정은 전략 코드에서 정밀 계산 (체결강도·분봉 반등 등) 하면
충분. 유니버스 단계에선 **활발한 종목 풀** 만 뽑으면 됨.
* 거래량/체결강도 상위는 KIS REST 한 방에 100~200건을 받아올 수 있고
재현 가능(tick 단위 저장 가능) → 백테스트 친화.
동작:
1. configs 리스트로 전략별 랭킹 소스 지정.
configs = [
{"strategy_id": "SCALP", "sort": "volume", "limit": 100, "market": "J"},
{"strategy_id": "SHORT", "sort": "volume", "limit": 100, "market": "J"},
# sort 후보: volume | trading_value | strength | fluct_up | fluct_down
]
2. 동일 (sort, market, limit) 조합은 REST 1회만 호출 (캐시).
3. 변동 감지 tick 마다 target_candidates_history 에 초단위 ``event_time`` 으로
**즉시 INSERT** (ConditionSearchManager 와 동일 패턴).
— 버퍼/배치 flush 는 단일 행 INSERT 대비 실익이 없고 crash 내구성만 낮추므로
제거. INSERT 1회 ≈ 1ms 미만이라 10초 폴링 루프에 영향 없음.
ConditionSearchManager 와의 공용 인터페이스 (전략 쪽에서는 구분 불필요):
- start() / stop()
- get_universe_for(strategy_id) -> Set[str]
- get_candidates_for(strategy_id) -> List[Dict]
- get_universe() -> Set[str] (전체 합집합)
- _configs 속성 (heartbeat 용)
"""
from __future__ import annotations
import random
import threading
import time
from datetime import datetime as dt
from typing import Callable, Dict, List, Optional, Set, Tuple
from ..utils.env import get_env_bool, get_env_int
from ..utils.logger import get_logger
logger = get_logger("kis_trader.rank")
_SORT_MAP = {
# 공개 별칭 → KIS FID_BLNG_CLS_CODE 매핑
"volume": "0",
"vol": "0",
"trading_value": "3",
"value": "3",
"strength": "6",
"cntr_str": "6",
"fluct_up": "4",
"up": "4",
"fluct_down": "5",
"down": "5",
"decline": "5",
}
class VolumeRankManager:
"""거래량/체결강도 순위 기반 동적 유니버스 매니저.
저장 정책: 변동 감지 tick 마다 초단위 ``event_time`` 으로 DB 에 **즉시 INSERT**.
(과거엔 큐 + 기록원 스레드로 배치 flush 했으나, 단일 행 INSERT 대비 실익이
없고 crash 시 미-flush tick 유실 위험만 낮추므로 제거했다.)
"""
def __init__(
self,
*,
client,
configs: List[Dict],
on_change: Optional[Callable[[str, Set[str], Set[str], Set[str]], None]] = None,
poll_interval_sec: Optional[float] = None,
db=None,
):
self.client = client
self.on_change = on_change
self.db = db
self.poll_interval = float(
poll_interval_sec
if poll_interval_sec is not None
else get_env_int("RANKING_POLL_INTERVAL_SEC", 10)
)
# 히스토리 저장 플래그 (ConditionSearchManager 와 env 공유 의도로 이름은 UNIVERSE_HISTORY_SAVE)
self.history_enabled = get_env_bool("UNIVERSE_HISTORY_SAVE", True)
# 설정 정규화
self._configs: List[Dict] = []
for c in configs or []:
sid = str(c.get("strategy_id") or "").strip().upper()
if not sid:
continue
sort_key = str(c.get("sort") or "volume").strip().lower()
blng = _SORT_MAP.get(sort_key)
if not blng:
logger.warning(
"⚠️ 알 수 없는 sort=%s (strategy=%s) → 'volume' 사용",
sort_key, sid,
)
sort_key = "volume"
blng = "0"
self._configs.append({
"strategy_id": sid,
"sort": sort_key,
"blng": blng,
"market": str(c.get("market") or "J").strip() or "J",
"limit": int(c.get("limit") or 100),
"exclude_non_stock": bool(c.get("exclude_non_stock", True)),
})
self._thread: Optional[threading.Thread] = None
self._running = False
self._current: Dict[str, Set[str]] = {}
self._name_map: Dict[str, str] = {}
self._initialized: Set[str] = set()
self._lock = threading.Lock()
# ------------------------------------------------------------------
# Public API
# ------------------------------------------------------------------
def start(self) -> bool:
if not self._configs:
logger.info("랭킹 매니저 configs 없음 → 비활성")
return False
for cfg in self._configs:
logger.info(
"🔗 랭킹 매핑: strategy=%s sort=%s market=%s limit=%d",
cfg["strategy_id"], cfg["sort"], cfg["market"], cfg["limit"],
)
self._running = True
self._thread = threading.Thread(
target=self._loop, daemon=True, name="VolumeRank",
)
self._thread.start()
logger.info(
"✅ 랭킹 폴링 시작 (%d개, interval=%ds, history=%s)",
len(self._configs), int(self.poll_interval),
"ON" if (self.history_enabled and self.db is not None) else "OFF",
)
return True
def stop(self) -> None:
"""폴링 스레드 정리."""
self._running = False
def get_universe_for(self, strategy_id: str) -> Set[str]:
sid = (strategy_id or "").upper()
with self._lock:
return set(self._current.get(sid, set()))
def get_candidates_for(self, strategy_id: str) -> List[Dict]:
"""BaseStrategy._load_candidates 와 호환되는 dict 리스트."""
sid = (strategy_id or "").upper()
with self._lock:
codes = list(self._current.get(sid, set()))
nm = dict(self._name_map)
out = []
for c in codes:
out.append({
"code": c,
"name": nm.get(c, c),
"scalp_on": True,
"tail_on": True,
"score": 0.0,
"price": 0.0,
})
return out
def get_universe(self) -> Set[str]:
"""전체 합집합 (heartbeat/총량 로그용)."""
with self._lock:
out: Set[str] = set()
for s in self._current.values():
out |= s
return out
def get_candidates(self) -> List[Dict]:
if not self._configs:
return []
return self.get_candidates_for(self._configs[0]["strategy_id"])
# ------------------------------------------------------------------
# 내부
# ------------------------------------------------------------------
def _loop(self) -> None:
while self._running:
try:
self._tick_all()
except Exception as e:
logger.error("랭킹 루프 예외: %s", e)
jitter = min(1.5, self.poll_interval * 0.1)
sleep_sec = self.poll_interval + random.uniform(0, jitter)
deadline = time.time() + sleep_sec
while time.time() < deadline:
if not self._running:
return
time.sleep(0.5)
def _tick_all(self) -> None:
# (blng, market, limit) 조합이 같으면 REST 1회만 호출
cache: Dict[Tuple[str, str, int, bool], List[Dict]] = {}
for cfg in self._configs:
key = (cfg["blng"], cfg["market"], cfg["limit"], cfg["exclude_non_stock"])
if key in cache:
rows = cache[key]
else:
try:
rows = self.client._fetch_volume_rank(
market=cfg["market"],
blng_cls_code=cfg["blng"],
limit=cfg["limit"],
exclude_non_stock=cfg["exclude_non_stock"],
) or []
except Exception as e:
logger.debug(
"랭킹 조회 실패 (blng=%s market=%s): %s",
cfg["blng"], cfg["market"], e,
)
rows = []
cache[key] = rows
self._apply_result(cfg["strategy_id"], rows)
def _apply_result(self, strategy_id: str, rows: List[Dict]) -> None:
new_set: Set[str] = set()
new_names: Dict[str, str] = {}
for r in rows:
code = (
r.get("mksc_shrn_iscd") or r.get("stk_cd")
or r.get("code") or ""
).strip()
if not code or len(code) != 6:
continue
name = (
r.get("hts_kor_isnm") or r.get("stk_nm")
or r.get("prst_name") or code
).strip() or code
new_set.add(code)
new_names[code] = name
with self._lock:
prev = self._current.get(strategy_id, set())
first_tick = strategy_id not in self._initialized
enters = new_set - prev
exits = prev - new_set
self._current[strategy_id] = new_set
self._initialized.add(strategy_id)
for c, n in new_names.items():
self._name_map[c] = n
changed = bool(enters or exits)
if changed:
logger.info(
"🔄 [%s] +%d / -%d (현재 %d종목, rank)",
strategy_id, len(enters), len(exits), len(new_set),
)
if enters:
preview = ", ".join(
f"{c}({new_names.get(c, c)})" for c in list(sorted(enters))[:5]
)
logger.info(" ENTER: %s%s",
preview, "" if len(enters) > 5 else "")
if exits:
preview = ", ".join(sorted(exits)[:5])
logger.info(" EXIT : %s%s",
preview, "" if len(exits) > 5 else "")
if self.on_change and changed:
try:
self.on_change(strategy_id, new_set, enters, exits)
except Exception as e:
logger.warning("on_change 콜백 예외: %s", e)
# 변동 또는 첫 tick → 풀 스냅샷 저장 (condition 과 동일 테이블/포맷)
if changed or first_tick:
self._save_snapshot(strategy_id, new_set, new_names)
def _save_snapshot(
self, strategy_id: str, codes: Set[str], names: Dict[str, str],
) -> None:
"""변동 감지 tick 마다 즉시 DB 에 풀 스냅샷 INSERT.
* ``event_time`` 은 REST 응답 직후 찍은 초단위 시각 (YYYY-MM-DD HH:MM:SS).
* 이 값이 백테스트 매칭 기준 시각이며, DB INSERT 가 1ms 지연되든 100ms
지연되든 행에 박히는 event_time 은 불변이므로 정합성은 보존.
* 폴링은 10초 주기라 DB 부하 매우 적음. 버퍼링으로 얻을 이득 없음.
"""
if not (self.history_enabled and self.db is not None):
return
event_time = dt.now().strftime("%Y-%m-%d %H:%M:%S")
items = [{"code": c, "name": names.get(c, c)} for c in sorted(codes)]
try:
n = self.db.insert_condition_universe_snapshot(
strategy_id=strategy_id,
event_time=event_time,
items=items,
)
logger.debug(
"📼 [history/rank] %s @%s %d종목 INSERT",
strategy_id, event_time, n,
)
except Exception as e:
# INSERT 실패는 매매 흐름에 영향 없어야 함 → 경고만 남기고 계속.
logger.warning("유니버스 히스토리 INSERT 실패 (%s): %s", strategy_id, e)

View File

@@ -0,0 +1,524 @@
"""
kis_trader/network/ws_manager.py — 단일 WebSocket 허브 (Event Bus)
====================================================================
설계 목적:
* 두 전략(스캘핑/꼬리잡기)이 각자 WS 연결을 띄우면 같은 종목 2번 구독 → 토큰/approval_key
경합 + 계정 차단 위험 → 프로세스 전체에서 **단 하나의 WS 세션**만 띄운다.
* 기존 ``kis_ws.KISWebSocketPriceCache`` + ``CandleAggregator`` 를 그대로 재사용.
* 전략별 "구독 관심 종목" 을 **레퍼런스 카운팅**으로 관리. 한 전략이 구독 해제해도
다른 전략이 구독 중이면 WS 에서 해제되지 않는다.
공용 API:
- start() / stop()
- subscribe(code, owner) / unsubscribe(code, owner)
- sync_targets(owner, codes) : 한 전략의 관심 종목을 한 번에 동기화
- get_price(code) / get_candles(code, tf, n)
- fill_gap(codes=None) : 갭 보정 (WS 연결 직후 자동 + 수동 호출)
"""
from __future__ import annotations
import queue
import random
import threading
import time
from collections import defaultdict
from typing import Dict, Iterable, List, Optional, Set
from ..utils.env import get_env_bool, get_env_from_db, get_env_int
from ..utils.logger import get_logger
logger = get_logger("kis_trader.ws")
# 기존 kis_ws 모듈 재사용 (검증된 로직 보존 원칙)
# 역할 분리 정책 (kis_scalping_ver2 / kis_short_ver3 와 동일):
# - 실시간 시세(WS) : KIS 실전키 (is_mock=False 고정)
# - 매수/매도 주문·계좌 : KIS (mock 여부는 KIS_MOCK)
# - 유니버스(거래량 순위) : KIS REST volume-rank
# - 과거 봉 웜업·갭보정 : 키움 ka10080 (1/3/15/60분 native 지원)
# → 키움 키 없으면 KIS 1분봉 fallback
try:
from kis_ws import (
CandleAggregator,
KISWebSocketPriceCache,
_get_kiwoom_creds,
get_kiwoom_candles_df,
)
_KIS_WS_AVAILABLE = True
except ImportError as _e:
_KIS_WS_AVAILABLE = False
# 서브심볼 import 실패 대응 (패키지만 있고 키움 함수 없는 구버전)
try:
from kis_ws import CandleAggregator, KISWebSocketPriceCache # type: ignore
_KIS_WS_AVAILABLE = True
_get_kiwoom_creds = None # type: ignore[assignment]
get_kiwoom_candles_df = None # type: ignore[assignment]
logger.warning(
"kis_ws 에 키움 함수 없음 → 갭보정 키움 fallback 비활성 "
"(KIS 1분봉 전용): %s", _e,
)
except ImportError as _e2:
logger.warning("kis_ws 모듈 import 실패 → WS 기능 비활성: %s", _e2)
class WSManager:
"""
단일 WS 허브. 전략은 이 매니저를 공유하고, subscribe/unsubscribe 시 owner 를 전달한다.
reference counting 예시:
subscribe("005930", owner="SCALP") # refs[005930]={SCALP} → WS subscribe
subscribe("005930", owner="SHORT") # refs[005930]={SCALP,SHORT} → (already subscribed)
unsubscribe("005930", owner="SCALP") # refs[005930]={SHORT} → keep
unsubscribe("005930", owner="SHORT") # refs[005930]=set() → WS unsubscribe
"""
def __init__(self, *, db, kis_client):
self.db = db
self.kis_client = kis_client
self.ws_cache: Optional["KISWebSocketPriceCache"] = None
self.candle_agg: Optional["CandleAggregator"] = None
# owner(전략ID) → 관심 코드 집합
self._owner_codes: Dict[str, Set[str]] = defaultdict(set)
# code → 보유 중인 owner 집합 (ref counting)
self._code_refs: Dict[str, Set[str]] = defaultdict(set)
# 영구 구독(시장방향 ETF 등)
self._permanent_codes: Set[str] = set()
self._lock = threading.Lock()
# ── 갭보정 비동기 파이프라인 ─────────────────────────────
# (전략 쓰레드에서 subscribe() 시 동기 REST 호출하면 매수 체크가
# 수 분간 블로킹됨 → 백그라운드 워커 큐로 이관)
self._gap_q: "queue.Queue[str]" = queue.Queue(maxsize=1024)
self._gap_filled: Set[str] = set() # 이미 갭보정 완료한 코드
self._gap_inflight: Set[str] = set() # 큐에 등록/처리 중인 코드
self._gap_lock = threading.Lock()
self._gap_worker_thread: Optional[threading.Thread] = None
# 전체 재갭보정(재접속 시) 중복 트리거 방지
self._bulk_refill_running = False
# ------------------------------------------------------------------
# 시작/종료
# ------------------------------------------------------------------
def start(self) -> bool:
"""WS 세션 시작. 실패 시 False (봇은 REST 폴백으로 동작)."""
if not _KIS_WS_AVAILABLE:
logger.warning("kis_ws 미설치 → WS 허브 비활성 (REST 폴백만 동작)")
return False
# ── [중요] WS 는 데이터 수신용이므로 무조건 실전 서버로 접속 ──
# kis_scalping_ver2 와 동일 정책 (모의 계좌라도 시세는 실전 필요)
ws_app_key = get_env_from_db("KIS_APP_KEY_REAL", "") or self.kis_client.app_key
ws_app_secret = get_env_from_db("KIS_APP_SECRET_REAL", "") or self.kis_client.app_secret
if not ws_app_key or not ws_app_secret:
logger.warning("KIS 실전 키 없음 → WS 허브 비활성")
return False
try:
self.ws_cache = KISWebSocketPriceCache(
app_key=ws_app_key,
app_secret=ws_app_secret,
is_mock=False, # 시세는 실전 서버 고정
)
# 봉 타임프레임: 스캘핑(1분) + 꼬리잡기(3분) + 추세 필터(15/60분)
# SCALP/SHORT 양쪽 전략이 쓰는 모든 TF 포함
tfs = self._resolve_timeframes()
self.candle_agg = CandleAggregator(db=self.db, timeframes=tfs)
self.ws_cache.attach_candle_aggregator(self.candle_agg)
ok = self.ws_cache.start()
if not ok:
logger.warning("WS start() 실패 → 비활성")
self.ws_cache = None
self.candle_agg = None
return False
# 갭보정 백그라운드 워커 가동 (subscribe 논블로킹 보장)
self._start_gap_worker()
# 영구 구독 (KOSPI/KOSDAQ ETF 등)
self._load_permanent_codes()
for code in sorted(self._permanent_codes):
self.ws_cache.subscribe(code)
self._enqueue_gap_fill(code)
logger.info("📡 [영구구독] %s", code)
# 연결 성공 후 자동 갭 보정 등록 (WS 재접속 시 전체 재갭보정)
self.ws_cache.set_on_connected_callback(self._trigger_bulk_refill_async)
logger.info(
"✅ WSManager 활성 (tfs=%s, permanent=%d, gap_worker=ON)",
tfs, len(self._permanent_codes),
)
return True
except Exception as e:
logger.error("WS 초기화 예외: %s", e)
self.ws_cache = None
self.candle_agg = None
return False
def stop(self) -> None:
try:
if self.ws_cache:
self.ws_cache.stop(clear_subscriptions=True)
except Exception as e:
logger.debug("WS stop 실패: %s", e)
@property
def is_active(self) -> bool:
return bool(self.ws_cache and self.ws_cache.is_active)
# ------------------------------------------------------------------
# 구독 관리 (Reference Counting)
# ------------------------------------------------------------------
def subscribe(self, code: str, owner: str) -> None:
"""
한 전략(owner)이 해당 종목에 관심 등록.
※ 갭보정은 **백그라운드 워커 큐**로 위임하여 전략 쓰레드를 블록하지 않음.
(예전: 여기서 REST 4개 TF 동기 호출 → 매수 체크 2~3분 지연)
"""
if not code or not owner:
return
with self._lock:
self._owner_codes[owner].add(code)
first_ref = not self._code_refs[code]
self._code_refs[code].add(owner)
if first_ref and self.ws_cache:
self.ws_cache.subscribe(code)
# 신규 구독 → 워커에게 갭보정 위임 (논블로킹)
self._enqueue_gap_fill(code)
def unsubscribe(self, code: str, owner: str) -> None:
"""한 전략(owner)이 관심 해제. 다른 전략이 아직 들고 있으면 WS 는 유지."""
if not code or not owner:
return
with self._lock:
self._owner_codes[owner].discard(code)
if owner in self._code_refs.get(code, set()):
self._code_refs[code].discard(owner)
still_refs = bool(self._code_refs.get(code))
is_permanent = code in self._permanent_codes
if not still_refs and not is_permanent and self.ws_cache:
self.ws_cache.unsubscribe(code)
if self.candle_agg:
self.candle_agg.remove_code(code)
def sync_targets(self, owner: str, codes: Iterable[str]) -> None:
"""
한 전략의 관심 종목 목록을 통째로 동기화.
- 기존 관심 종목 중 없어진 것은 unsubscribe
- 새로 추가된 것은 subscribe
"""
new_set = {c for c in codes if c}
with self._lock:
cur = set(self._owner_codes.get(owner, set()))
for code in sorted(cur - new_set):
self.unsubscribe(code, owner)
for code in sorted(new_set - cur):
self.subscribe(code, owner)
# ------------------------------------------------------------------
# 조회 헬퍼 (전략이 쓰는 API)
# ------------------------------------------------------------------
def get_price(self, code: str, max_age_sec: float = 5.0) -> Optional[dict]:
if not self.ws_cache:
return None
try:
return self.ws_cache.get_price(code, max_age_sec=max_age_sec)
except Exception:
return None
def get_candles(self, code: str, tf: int, n: int = 50) -> list:
if self.candle_agg:
try:
return self.candle_agg.get_candles(code, tf, n)
except Exception:
return []
# CandleAggregator 없으면 DB 폴백
try:
return self.db.get_ws_candles(code, tf, limit=n, confirmed_only=True)
except Exception:
return []
def fill_gap(self, codes: Optional[Iterable[str]] = None) -> None:
"""외부에서 수동으로 갭 보정 트리거 (비동기: 큐 등록 후 즉시 리턴)."""
if codes is None:
self._trigger_bulk_refill_async()
else:
for c in codes:
self._enqueue_gap_fill(c)
# ------------------------------------------------------------------
# 내부: 갭 보정 — 백그라운드 워커 파이프라인
# ------------------------------------------------------------------
def _start_gap_worker(self) -> None:
"""갭보정 전담 데몬 워커 스레드 기동 (단일 워커 → REST 레이트리밋 자연 직렬화)."""
if self._gap_worker_thread and self._gap_worker_thread.is_alive():
return
t = threading.Thread(
target=self._gap_worker_loop,
name="WS-GapFillWorker",
daemon=True,
)
t.start()
self._gap_worker_thread = t
logger.info("✅ 갭보정 워커 스레드 시작 (queue 기반 비동기 처리)")
def _enqueue_gap_fill(self, code: str) -> None:
"""구독 직후 호출 — 갭보정 큐에 논블로킹 등록.
중복 방지:
- 이미 완료(`_gap_filled`) → 스킵
- 이미 큐/처리 중(`_gap_inflight`) → 스킵
"""
if not code:
return
with self._gap_lock:
if code in self._gap_filled or code in self._gap_inflight:
return
self._gap_inflight.add(code)
try:
self._gap_q.put_nowait(code)
except queue.Full:
# 큐가 가득 차면 inflight 해제 후 포기 (WS 틱으로 자연 누적)
with self._gap_lock:
self._gap_inflight.discard(code)
logger.warning("⚠️ 갭보정 큐 full → %s 스킵 (WS 실시간 누적으로 대체)", code)
def _trigger_bulk_refill_async(self) -> None:
"""WS 재접속 시 현재 구독된 전 종목의 갭보정 완료 마커를 리셋하고 재큐잉."""
if not (self.ws_cache and self.candle_agg):
return
if self._bulk_refill_running:
return
self._bulk_refill_running = True
def _bulk():
try:
with self.ws_cache._sub_lock:
codes = sorted(self.ws_cache._subscribed)
# 재접속이므로 모든 종목 갭보정 재실행
with self._gap_lock:
self._gap_filled.clear()
logger.info(
"🔄 [갭보정-전체] WS 재접속 → %d종목 큐 재등록", len(codes),
)
for code in codes:
self._enqueue_gap_fill(code)
finally:
self._bulk_refill_running = False
threading.Thread(target=_bulk, name="WS-BulkRefill", daemon=True).start()
def _gap_worker_loop(self) -> None:
"""단일 워커 루프: 큐에서 코드 꺼내 순차 처리 → REST 레이트리밋 자연 완충."""
# 크레덴셜은 첫 작업 시점에 1회 조회 후 캐시 (env 변경 무시하고 세션 유지)
kw_key = kw_secret = None
kw_mock = False
kw_resolved = False
while True:
try:
code = self._gap_q.get(timeout=1.0)
except queue.Empty:
continue
if code is None: # 종료 시그널
return
# 장중만 실행 (장외면 완료 마커 찍고 다음)
if not self._is_market_hours() and not get_env_bool("WS_GAP_FILL_OFF_HOURS", False):
with self._gap_lock:
self._gap_inflight.discard(code)
self._gap_filled.add(code)
self._gap_q.task_done()
continue
if not kw_resolved:
kw_key, kw_secret, kw_mock = self._get_kiwoom_credentials()
use_kiwoom = bool(kw_key and kw_secret and get_kiwoom_candles_df is not None)
if use_kiwoom:
kw_status = f"✅ ({'모의' if kw_mock else '실전'})"
else:
kw_status = ""
logger.info(
"🔧 [갭보정-워커] kiwoom=%s, KIS_fallback=%s",
kw_status,
"ON" if get_env_bool("WS_GAP_FILL_KIS_FALLBACK", False) else "OFF",
)
kw_resolved = True
try:
self._fill_gap_for_code(
code,
kw_key=kw_key, kw_secret=kw_secret, kw_mock=kw_mock,
)
except Exception as e:
logger.debug("갭보정 워커 예외 (%s): %s", code, e)
finally:
with self._gap_lock:
self._gap_inflight.discard(code)
self._gap_filled.add(code)
self._gap_q.task_done()
# 종목 간 짧은 sleep (REST 레이트리밋 완충)
time.sleep(random.uniform(0.2, 0.4))
@staticmethod
def _is_market_hours() -> bool:
"""
09:00~15:30 KST 평일 여부.
KIS 모의투자 서버는 장외시간 inquire-time-itemchartprice 호출에
HTTP 500 을 던지므로, 갭보정 REST 호출은 장중에만 시도한다.
(키움 ka10080 은 장외에도 동작하지만, 매매 자체가 장중에만 의미 있으므로 통일)
"""
now = time.localtime()
if now.tm_wday >= 5: # 토/일
return False
hhmm = now.tm_hour * 100 + now.tm_min
return 900 <= hhmm <= 1530
def _get_kiwoom_credentials(self):
"""
키움 분봉 갭보정용 키 조회 — **KIS 처럼 모의/실전 토글 가능**.
토글 결정 우선순위
------------------
1) ``KIWOOM_MOCK`` (있으면 단독 사용 — 키움만 별도 토글하고 싶을 때)
2) 미지정 시 ``KIS_MOCK`` 폴백 (한 번만 설정해도 동기화)
키 슬롯 매핑
-----------
mock=True → ``KIWOOM_APP_KEY_MOCK`` → 없으면 ``KIWOOM_APP_KEY`` (레거시) 폴백
mock=False → ``KIWOOM_APP_KEY_REAL`` → 없으면 ``KIWOOM_APP_KEY`` (레거시) 폴백
키 없으면 ``(None, None, is_mock)`` 반환 → 키움 비활성, KIS REST fallback 사용
(``WS_GAP_FILL_KIS_FALLBACK`` 권장 ON)
주의
----
레거시 폴백은 키-도메인이 어긋나면 키움이 ``8030`` 으로 거부한다.
(예: 모의 전용 키를 실전 도메인에 던지면 8030.) 폴백 사용 시 로그
한 줄로 명시한다.
Returns:
(app_key, app_secret, is_mock)
"""
if get_kiwoom_candles_df is None:
return None, None, False
try:
# ── 1. 토글 결정 ──────────────────────────────────────────
kw_mock_raw = (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 = get_env_bool("KIS_MOCK", True)
# ── 2. 키 슬롯 선택 (모의/실전) ───────────────────────────
if is_mock:
kw_key = (get_env_from_db("KIWOOM_APP_KEY_MOCK", "") or "").strip()
kw_secret = (get_env_from_db("KIWOOM_APP_SECRET_MOCK", "") or "").strip()
else:
kw_key = (get_env_from_db("KIWOOM_APP_KEY_REAL", "") or "").strip()
kw_secret = (get_env_from_db("KIWOOM_APP_SECRET_REAL", "") or "").strip()
# ── 3. 레거시 단일 필드 폴백 (KIWOOM_APP_KEY/_SECRET) ────
if not kw_key or not kw_secret:
legacy_key = (get_env_from_db("KIWOOM_APP_KEY", "") or "").strip()
legacy_secret = (get_env_from_db("KIWOOM_APP_SECRET", "") or "").strip()
if legacy_key and legacy_secret:
kw_key = kw_key or legacy_key
kw_secret = kw_secret or legacy_secret
logger.info(
"🔧 [키움] %s 슬롯 비어있어 LEGACY KIWOOM_APP_KEY 폴백 사용 "
"(키-도메인 불일치 시 8030 발생 가능)",
"MOCK" if is_mock else "REAL",
)
if not kw_key or not kw_secret:
return None, None, is_mock
return kw_key, kw_secret, is_mock
except Exception as e:
logger.debug("키움 크레덴셜 조회 예외: %s", e)
return None, None, False
def _fill_gap_for_code(
self,
code: str,
*,
kw_key: Optional[str] = None,
kw_secret: Optional[str] = None,
kw_mock: bool = False,
) -> None:
"""
단일 종목 갭 보정 — 워커 스레드 전용 (전략 쓰레드에서 직접 호출 금지).
정책 (기존 kis_scalping_ver2._fill_all_gaps 개선):
[1] 키움 ka10080 우선 (1/3/15/60분봉 native + 과거봉 확보)
[2] 키움 실패 시 KIS fallback — **기본 OFF** (`WS_GAP_FILL_KIS_FALLBACK=false`)
→ KIS 모의 서버가 장중에도 HTTP 500 을 자주 반환해 로그 오염 + 백오프 지연.
키움 있으면 굳이 안 쳐도 됨. WS 틱이 쌓여 자연 보완됨.
→ env 로 true 지정 시에만 tf<=3 한정 KIS 호출.
[3] 어느 경로든 실패 → WS 실시간 틱으로 자연 누적 (CandleAggregator)
"""
if not (self.ws_cache and self.candle_agg):
return
if not self._is_market_hours() and not get_env_bool("WS_GAP_FILL_OFF_HOURS", False):
return
limit = get_env_int("WS_GAP_FILL_LIMIT", 120)
use_kiwoom = bool(kw_key and kw_secret and get_kiwoom_candles_df is not None)
kis_fallback_on = get_env_bool("WS_GAP_FILL_KIS_FALLBACK", False)
for tf in self.candle_agg.timeframes:
df = None
if use_kiwoom:
try:
df = get_kiwoom_candles_df(
code, tf, kw_key, kw_secret,
is_mock=kw_mock, n=limit,
)
except Exception as e:
logger.debug("키움 갭보정 실패 (%s %dM): %s", code, tf, e)
# KIS fallback — env 로 명시적 ON 일 때만 (1/3분봉 한정)
if (df is None or df.empty) and kis_fallback_on and tf <= 3:
try:
df = self.kis_client.get_minute_chart(
code, period=str(tf), limit=limit,
)
except Exception as e:
logger.debug("KIS 갭보정 실패 (%s %dM): %s", code, tf, e)
if df is not None and not df.empty:
self.candle_agg.fill_gap_from_rest(code, tf, df)
# 같은 종목 내 타임프레임 전환 사이 짧은 sleep (차트 API 레이트리밋)
time.sleep(random.uniform(0.15, 0.3))
# ------------------------------------------------------------------
# 영구 구독 / 타임프레임 해석
# ------------------------------------------------------------------
def _load_permanent_codes(self) -> None:
raw = get_env_from_db("PERMANENT_WS_CODES", "069500,229200")
self._permanent_codes = {c.strip() for c in str(raw).split(",") if c.strip()}
def _resolve_timeframes(self) -> List[int]:
"""SCALP 1분 + SHORT 3분 + 추세 15/60분. env 로 확장 가능."""
tf_raw = get_env_from_db("WS_TIMEFRAMES", "1,3,15,60")
try:
tfs = [int(x.strip()) for x in str(tf_raw).split(",") if x.strip()]
# 최소 1분, 3분은 포함 보장 (두 전략 필수 TF)
for must in (1, 3):
if must not in tfs:
tfs.append(must)
return sorted(set(tfs))
except Exception:
return [1, 3, 15, 60]