""" kis_trader/ws/orderbook_cache.py — 키움 0D 호가잔량 RAM 캐시 ================================================================ 키움 OpenAPI+ FID (주식호가잔량 0D): 매도호가 1~10: 41~50, 매수호가 1~10: 51~60 매도수량 1~10: 61~70, 매수수량 1~10: 71~80 매도총잔량: 121, 매수총잔량: 125 RAM 적재: 틱과 동일 — 거래소 시각(snap_time)이 LIVE_FEED_FALLBACK(기본 3초)보다 오래면 넣지 않음(ts만 새로 찍어 낡은 호가가 3초 체인을 통과하지 못하게). """ from __future__ import annotations import json import logging import threading import time from dataclasses import dataclass, field from typing import Any, Dict, List, Optional, Tuple logger = logging.getLogger("kis_trader.orderbook_cache") def _abs_int(v: Any) -> int: try: return int(abs(float(str(v or "0").replace(",", "")))) except (TypeError, ValueError): return 0 def _normalize_ob_snap_time(raw: Any) -> str: """FID20 / hotime / BSOP_HOUR → YYYYMMDDHHMMSS (가능하면).""" tt = str(raw or "").strip().replace(":", "").replace("-", "").replace(" ", "") if not tt: return "" if len(tt) >= 14 and tt[:14].isdigit(): return tt[:14] if len(tt) >= 6 and tt[-6:].isdigit(): from datetime import datetime as _dt return _dt.now().strftime("%Y%m%d") + tt[-6:] return "" def _is_snap_time_stale_for_ram(snap: "OrderbookSnapshot") -> bool: """틱 skip_ram 과 동일 — LIVE_FEED_FALLBACK 기준 snap_time lag.""" try: from kis_trader.engine.feed_fallback import is_orderbook_snap_time_stale return bool(is_orderbook_snap_time_stale(getattr(snap, "snap_time", "") or "")) except Exception: return False @dataclass class OrderbookLevel: price: int = 0 qty: int = 0 @dataclass class OrderbookSnapshot: """종목 1개 호가 스냅샷 (매수호가=비싼 순, 매도호가=싼 순).""" code: str bids: List[OrderbookLevel] = field(default_factory=list) asks: List[OrderbookLevel] = field(default_factory=list) total_bid_qty: int = 0 total_ask_qty: int = 0 ts: float = 0.0 source: str = "kiwoom_0d" snap_time: str = "" # YYYYMMDDHHMMSS — DB·판정 스냅 재생용 def best_bid(self) -> int: return self.bids[0].price if self.bids else 0 def best_ask(self) -> int: return self.asks[0].price if self.asks else 0 def spread_pct(self) -> float: bb, ba = self.best_bid(), self.best_ask() if bb <= 0 or ba <= 0: return 999.0 mid = (bb + ba) / 2.0 if mid <= 0: return 999.0 return (ba - bb) / mid * 100.0 def bid_qty_sum(self, levels: int = 3) -> int: n = max(1, int(levels)) return sum(lv.qty for lv in self.bids[:n]) def ask_qty_sum(self, levels: int = 3) -> int: n = max(1, int(levels)) return sum(lv.qty for lv in self.asks[:n]) def to_kis_bid_dict(self, levels: int = 10) -> Dict[str, Any]: """``orderbook_sell.parse_kis_bid_levels`` 호환 dict.""" out: Dict[str, Any] = {} for i, lv in enumerate(self.bids[: max(1, levels)], start=1): out[f"bidp{i}"] = lv.price out[f"bidp_rsqn{i}"] = lv.qty return out def to_storage_dict(self) -> Dict[str, Any]: """``ws_orderbook`` INSERT·백테 복원용.""" return { "code": self.code, "best_bid": self.best_bid(), "best_ask": self.best_ask(), "total_bid_qty": self.total_bid_qty, "total_ask_qty": self.total_ask_qty, "bid_qty_l3": self.bid_qty_sum(3), "ask_qty_l3": self.ask_qty_sum(3), "levels_json": json.dumps( { "bids": [[lv.price, lv.qty] for lv in self.bids[:10]], "asks": [[lv.price, lv.qty] for lv in self.asks[:10]], }, ensure_ascii=False, separators=(",", ":"), ), "source": self.source, "ts": self.ts, "snap_time": self.snap_time or "", } def orderbook_snapshot_from_storage(row: Dict[str, Any]) -> OrderbookSnapshot: """DB 행 또는 storage dict → ``OrderbookSnapshot``.""" bids: List[OrderbookLevel] = [] asks: List[OrderbookLevel] = [] raw = row.get("levels_json") if raw: try: lv = json.loads(raw) if isinstance(raw, str) else raw for px, qty in (lv.get("bids") or []): bids.append(OrderbookLevel(price=_abs_int(px), qty=_abs_int(qty))) for px, qty in (lv.get("asks") or []): asks.append(OrderbookLevel(price=_abs_int(px), qty=_abs_int(qty))) except (TypeError, ValueError, json.JSONDecodeError): pass total_bid = _abs_int(row.get("total_bid_qty")) total_ask = _abs_int(row.get("total_ask_qty")) if total_bid <= 0 and bids: total_bid = sum(b.qty for b in bids) if total_ask <= 0 and asks: total_ask = sum(a.qty for a in asks) ts_raw = row.get("ts") try: ts = float(ts_raw) if ts_raw not in (None, "") else time.time() except (TypeError, ValueError): ts = time.time() st = str(row.get("snap_time") or "").strip() return OrderbookSnapshot( code=str(row.get("code") or "").strip(), bids=bids, asks=asks, total_bid_qty=total_bid, total_ask_qty=total_ask, ts=ts, source=str(row.get("source") or "kiwoom_0d"), snap_time=st[:14] if st else "", ) def parse_kiwoom_0d_values(code: str, values: Dict[str, Any]) -> OrderbookSnapshot: """키움 WS 0D ``values`` dict → ``OrderbookSnapshot``.""" asks: List[OrderbookLevel] = [] bids: List[OrderbookLevel] = [] for i in range(1, 11): ask_px = _abs_int(values.get(str(40 + i))) bid_px = _abs_int(values.get(str(50 + i))) ask_qty = _abs_int(values.get(str(60 + i))) bid_qty = _abs_int(values.get(str(70 + i))) if ask_px > 0: asks.append(OrderbookLevel(price=ask_px, qty=ask_qty)) if bid_px > 0: bids.append(OrderbookLevel(price=bid_px, qty=bid_qty)) total_ask = _abs_int(values.get("121")) total_bid = _abs_int(values.get("125")) if total_bid <= 0 and bids: total_bid = sum(b.qty for b in bids) if total_ask <= 0 and asks: total_ask = sum(a.qty for a in asks) # FID20 체결시각(HHMMSS) — 0D에도 실리면 틱과 같은 snap 축. 없으면 "" → lag 컷 스킵. snap_time = _normalize_ob_snap_time(values.get("20") or values.get(20)) return OrderbookSnapshot( code=str(code).strip(), asks=asks, bids=bids, total_bid_qty=total_bid, total_ask_qty=total_ask, ts=time.time(), source="kiwoom_0d", snap_time=snap_time, ) def _ls_level_qty(body: Dict[str, Any], side: str, i: int) -> int: """LS UH1/H1_/HA_ 잔량 — 통합(UH1)은 krx+nxt+unt 합산, 단일시장은 offerremN.""" # H1_/HA_: offerrem1 / bidrem1 plain = _abs_int(body.get(f"{side}rem{i}")) if plain > 0: return plain # UH1: krx_offerrem1 + nxt_ + unt_ return ( _abs_int(body.get(f"krx_{side}rem{i}")) + _abs_int(body.get(f"nxt_{side}rem{i}")) + _abs_int(body.get(f"unt_{side}rem{i}")) ) def _ls_total_qty(body: Dict[str, Any], side: str) -> int: """totofferrem / totoffer — UH1 은 krx_tot* + nxt_ + unt_ 합산.""" plain = _abs_int(body.get(f"tot{side}rem")) if plain > 0: return plain return ( _abs_int(body.get(f"krx_tot{side}rem")) + _abs_int(body.get(f"nxt_tot{side}rem")) + _abs_int(body.get(f"unt_tot{side}rem")) ) def parse_kis_h0stasp0_fields(code: str, fields: List[Any]) -> OrderbookSnapshot: """한투 실시간 호가 H0STASP0 ``^`` 필드 → ``OrderbookSnapshot``. 공식 columns (asking_price_krx / MCP): 0 MKSC_SHRN_ISCD, 1 BSOP_HOUR, 2 HOUR_CLS_CODE, 3~12 ASKP1~10, 13~22 BIDP1~10, 23~32 ASKP_RSQN1~10, 33~42 BIDP_RSQN1~10, 43 TOTAL_ASKP_RSQN, 44 TOTAL_BIDP_RSQN """ asks: List[OrderbookLevel] = [] bids: List[OrderbookLevel] = [] n = len(fields or []) for i in range(10): ask_px = _abs_int(fields[3 + i]) if n > 3 + i else 0 bid_px = _abs_int(fields[13 + i]) if n > 13 + i else 0 ask_qty = _abs_int(fields[23 + i]) if n > 23 + i else 0 bid_qty = _abs_int(fields[33 + i]) if n > 33 + i else 0 if ask_px > 0: asks.append(OrderbookLevel(price=ask_px, qty=ask_qty)) if bid_px > 0: bids.append(OrderbookLevel(price=bid_px, qty=bid_qty)) total_ask = _abs_int(fields[43]) if n > 43 else 0 total_bid = _abs_int(fields[44]) if n > 44 else 0 if total_bid <= 0 and bids: total_bid = sum(b.qty for b in bids) if total_ask <= 0 and asks: total_ask = sum(a.qty for a in asks) snap_time = _normalize_ob_snap_time(fields[1] if n > 1 else "") code_val = str(code or "").strip() if not code_val and n > 0: code_val = str(fields[0] or "").strip() return OrderbookSnapshot( code=code_val, asks=asks, bids=bids, total_bid_qty=total_bid, total_ask_qty=total_ask, ts=time.time(), source="kis_h0stasp0", snap_time=snap_time, ) def parse_ls_hoga_body( code: str, body: Dict[str, Any], *, source: str = "ls_uh1", ) -> OrderbookSnapshot: """LS 호가잔량 body (UH1/H1_/HA_) → ``OrderbookSnapshot``. - 가격: ``offerhoN``(매도)·``bidhoN``(매수) 1~10 - 잔량: 단일 ``offerremN`` 또는 통합 ``krx_/nxt_/unt_`` 합 """ asks: List[OrderbookLevel] = [] bids: List[OrderbookLevel] = [] for i in range(1, 11): ask_px = _abs_int(body.get(f"offerho{i}")) bid_px = _abs_int(body.get(f"bidho{i}")) ask_qty = _ls_level_qty(body, "offer", i) bid_qty = _ls_level_qty(body, "bid", i) if ask_px > 0: asks.append(OrderbookLevel(price=ask_px, qty=ask_qty)) if bid_px > 0: bids.append(OrderbookLevel(price=bid_px, qty=bid_qty)) total_ask = _ls_total_qty(body, "offer") total_bid = _ls_total_qty(body, "bid") if total_bid <= 0 and bids: total_bid = sum(b.qty for b in bids) if total_ask <= 0 and asks: total_ask = sum(a.qty for a in asks) # UH1 패킷 시각(HHMMSS). 틱동기 저장은 체결 chetime 으로 덮어씀(키움 FID20 과 동일 축). snap_time = _normalize_ob_snap_time(body.get("hotime") or "") return OrderbookSnapshot( code=str(code).strip(), asks=asks, bids=bids, total_bid_qty=total_bid, total_ask_qty=total_ask, ts=time.time(), source=str(source or "ls_uh1"), snap_time=snap_time, ) class OrderbookCache: """스레드 세이프 종목별 호가 캐시.""" def __init__(self) -> None: self._data: Dict[str, OrderbookSnapshot] = {} self._lock = threading.Lock() def _commit_if_fresh(self, snap: OrderbookSnapshot) -> Optional[OrderbookSnapshot]: """거래소 snap_time 이 LIVE_FEED_FALLBACK 초과면 RAM 미반영 (틱 skip_ram 과 동일). Returns: 저장된 snap. 스킵 시 None (호출측 recorder/DB 도 생략). """ if snap is None or not str(getattr(snap, "code", "") or "").strip(): return None if _is_snap_time_stale_for_ram(snap): logger.debug( "호가 RAM skip (snap_time lag) %s src=%s snap=%s", snap.code, getattr(snap, "source", ""), getattr(snap, "snap_time", "") or "-", ) return None with self._lock: self._data[snap.code] = snap return snap def update_from_kiwoom_0d( self, code: str, values: Dict[str, Any], ) -> Optional[OrderbookSnapshot]: snap = parse_kiwoom_0d_values(code, values) return self._commit_if_fresh(snap) def update_from_ls_hoga( self, code: str, body: Dict[str, Any], *, source: str = "ls_uh1", ) -> Optional[OrderbookSnapshot]: snap = parse_ls_hoga_body(code, body, source=source) return self._commit_if_fresh(snap) def update_from_kis_h0stasp0( self, code: str, fields: List[Any], ) -> Optional[OrderbookSnapshot]: snap = parse_kis_h0stasp0_fields(code, fields) return self._commit_if_fresh(snap) def get(self, code: str, max_age_sec: float = 3.0) -> Optional[OrderbookSnapshot]: c = str(code or "").strip() if not c: return None with self._lock: snap = self._data.get(c) if not snap: return None # 수신 ts 나이 (기존). max_age<=0 = 마지막 RAM(필터용) — ts·snap 컷 안 함. if max_age_sec > 0 and (time.time() - snap.ts) > max_age_sec: return None # 읽기 폴백 나이일 때 거래소 snap 도 틱과 동일 검사 (적재 누락·재시작 잔여 방어). # FILTER_MAX_AGE=0(마지막 RAM) 경로는 여기 안 탐 → 호가필터 구멍 방지 규칙 유지. if max_age_sec > 0 and _is_snap_time_stale_for_ram(snap): return None return snap def get_kis_bid_dict(self, code: str, max_age_sec: float = 3.0) -> Optional[Dict[str, Any]]: snap = self.get(code, max_age_sec=max_age_sec) if not snap or not snap.bids: return None return snap.to_kis_bid_dict() def remove(self, code: str) -> None: c = str(code or "").strip() with self._lock: self._data.pop(c, None)