""" kis_trader/engine/feed_fallback.py — 실매·옵투나 공통 읽기 나이 헬퍼 ====================================================================== 메인이 지금보다 LIVE_FEED_FALLBACK_MAX_AGE_SEC 초보다 오래면 그 벤더는 없는 것과 같음. 2차(보조 WS) → 3차(LS) → (실매 매도만) 4차 키움 REST. 적재(공책) 임계와 읽기 나이는 분리한다. OHLC 로 틱을 메우지 않음. """ from __future__ import annotations from datetime import datetime import time from typing import Any, Dict, List, Optional, Sequence, Tuple from kis_trader.utils.env import get_env_bool, get_env_float, get_env_from_db def packet_lag_seconds(pkt_raw: str, *, now_dt: Optional[datetime] = None) -> Optional[float]: """체결시각 문자열 vs now_dt(기본=서버 지금). 증권사끼리 비교 아님. YYYYMMDDHHMMSS 또는 HHMMSS. 파싱 실패면 None(나이를 모름 → 읽기에서 버리지 않음). """ tt = str(pkt_raw or "").strip() if not tt: return None wall = now_dt if now_dt is not None else datetime.now() pkt_dt = None try: if len(tt) >= 14 and tt[:14].isdigit(): pkt_dt = datetime.strptime(tt[:14], "%Y%m%d%H%M%S") elif len(tt) >= 6 and tt[-6:].isdigit(): pkt_dt = datetime.strptime( wall.strftime("%Y%m%d") + tt[-6:], "%Y%m%d%H%M%S", ) except Exception: return None if pkt_dt is None: return None return float((wall - pkt_dt).total_seconds()) def is_feed_read_stale( lag_sec: Optional[float], max_age_sec: Optional[float] = None, ) -> bool: """읽기 나이 초과. lag/나이를 모르면 False(버리지 않음).""" age = live_feed_fallback_max_age_sec() if max_age_sec is None else float(max_age_sec) if age <= 0 or lag_sec is None: return False try: return float(lag_sec) > age except (TypeError, ValueError): return False def is_orderbook_snap_time_stale( snap_time: Any, *, max_age_sec: Optional[float] = None, ) -> bool: """호가 거래소 시각(snap_time) vs 지금 — 틱 ``packet_lag`` + ``is_feed_read_stale`` 와 동일. snap_time 없거나 파싱 실패 → False(버리지 않음). LIVE_FEED_FALLBACK 0 이하면 OFF. """ return is_feed_read_stale( packet_lag_seconds(str(snap_time or "").strip()), max_age_sec=max_age_sec, ) def live_feed_fallback_max_age_sec() -> float: """읽기 유효 나이(초). 기본 3. 0 이하면 나이로 벤더를 버리지 않음(레거시).""" try: v = float(get_env_float("LIVE_FEED_FALLBACK_MAX_AGE_SEC", 3.0) or 3.0) except (TypeError, ValueError): v = 3.0 return v def live_tick_primary() -> str: p = str(get_env_from_db("LIVE_TICK_PROVIDER", "kiwoom") or "kiwoom").strip().lower() return p if p in ("kis", "kiwoom") else "kiwoom" def live_ob_primary() -> str: p = str(get_env_from_db("LIVE_OB_PROVIDER", "kiwoom") or "kiwoom").strip().lower() return p if p in ("kis", "kiwoom") else "kiwoom" def vendor_read_max_age_sec(caller_max_age: Optional[float]) -> Optional[float]: """벤더 1회 조회에 쓸 나이. - caller None = 마지막 RAM(명시, 나이 무시) - caller 0/음수 또는 생략 대체 = 폴백 나이(기본 3초) - caller 양수 = min(caller, 폴백나이). 폴백 0이면 caller 그대로. """ fb = live_feed_fallback_max_age_sec() if caller_max_age is None: return None try: c = float(caller_max_age) except (TypeError, ValueError): c = 0.0 if fb <= 0: return None if c <= 0 else c if c <= 0: return fb return min(c, fb) def tick_feed_tier(vendor: str, primary: str) -> int: v = str(vendor or "").strip().lower() p = str(primary or "kiwoom").strip().lower() if v in ("kiwoom_rest", "rest", "ka10007"): return 4 if v == "ls": return 3 if v == p: return 1 return 2 def format_vendor_label(vendor: str, tier: int = 0, spilled: bool = False) -> str: v = str(vendor or "").strip().lower() or "?" names = { "kiwoom": "kiwoom", "kis": "kis", "ls": "ls", "kiwoom_rest": "kiwoom_rest", "rest": "kiwoom_rest", "ka10007": "kiwoom_rest", } name = names.get(v, v) extra: List[str] = [] tmap = {1: "1차", 2: "2차", 3: "3차", 4: "4차"} tlab = tmap.get(int(tier or 0), "") if tlab: extra.append(tlab) if spilled: extra.append("spill") if extra: return f"{name}({','.join(extra)})" return name def format_mm_feed_line(tick_lab: str, ob_lab: str) -> str: t = str(tick_lab or "").strip() or "?" o = str(ob_lab or "").strip() or "?" return f"시세: {t} | 호가: {o}" def trigger_feed_detail_log_enabled() -> bool: """트리거(매수체크) 로그에 틱·호가 상세 추적 줄. 기본 ON.""" try: return bool(get_env_bool("TRIGGER_FEED_DETAIL_LOG", True)) except Exception: return True def extract_tick_trace_fields(data: Optional[Dict[str, Any]]) -> Dict[str, Any]: """get_price dict → 로그용 현재가·틱타임·나이.""" out: Dict[str, Any] = {} if not isinstance(data, dict): return out px = ( data.get("stck_prpr") or data.get("price") or data.get("cur_prc") or data.get("close") or 0 ) try: out["price"] = float(str(px).replace(",", "") or 0) except (TypeError, ValueError): out["price"] = 0.0 tt = ( data.get("tick_time") or data.get("tick_time_raw") or data.get("chetime") or data.get("cntr_tm") or data.get("FID20") or data.get("stck_cntg_hour") or data.get("cntg_hour") or "" ) out["tick_time"] = str(tt or "").strip() if not out["tick_time"]: bs = str(data.get("kis_bsop_date_raw") or "").strip() hr = str( data.get("kis_cntg_hour_raw") or data.get("chetime") or data.get("kiwoom_fid20") or "" ).strip() if bs and hr: out["tick_time"] = (bs + hr)[:14] elif hr: out["tick_time"] = hr # 키움 캐시에 tick_time 키가 없어도 chetime만으로 lag 계산 if data.get("chetime") and not out.get("tick_time"): out["tick_time"] = str(data.get("chetime") or "").strip() if data.get("_age_ms") is not None: try: out["age_ms"] = int(data.get("_age_ms") or 0) except (TypeError, ValueError): pass lag = packet_lag_seconds(out.get("tick_time") or "") if lag is not None: out["lag_sec"] = round(float(lag), 2) if data.get("_feed_vendor"): out["vendor"] = str(data.get("_feed_vendor") or "").strip().lower() return out def extract_ob_trace_fields(snap: Any) -> Dict[str, Any]: """호가 스냅샷 → bid/ask/스프레드/잔량비·스냅시각.""" out: Dict[str, Any] = {} if snap is None: return out try: if hasattr(snap, "best_bid"): out["best_bid"] = int(snap.best_bid() or 0) out["best_ask"] = int(snap.best_ask() or 0) try: out["spread_pct"] = round(float(snap.spread_pct() or 0), 3) except Exception: pass try: bq = int(getattr(snap, "total_bid_qty", 0) or 0) aq = int(getattr(snap, "total_ask_qty", 0) or 0) out["bid_qty"] = bq out["ask_qty"] = aq if aq > 0: out["bid_ask_ratio"] = round(bq / float(aq), 3) except Exception: pass out["snap_time"] = str(getattr(snap, "snap_time", "") or "").strip() src = str(getattr(snap, "source", "") or "").strip() if src: out["ob_source"] = src if getattr(snap, "ts", None): try: out["age_ms"] = int(max(0.0, (time.time() - float(snap.ts)) * 1000)) except Exception: pass return out except Exception: pass if isinstance(snap, dict): try: out["best_bid"] = int(float(snap.get("best_bid") or snap.get("bidp1") or 0)) out["best_ask"] = int(float(snap.get("best_ask") or snap.get("askp1") or 0)) except (TypeError, ValueError): pass out["snap_time"] = str(snap.get("snap_time") or "").strip() return out def format_trigger_feed_trace( tick_rec: Optional[Dict[str, Any]], ob_rec: Optional[Dict[str, Any]], *, tick_primary: str = "", ob_primary: str = "", ) -> str: """매수체크 로그 꼬리 — 벤더(1·2·3차) + 현재가 + 틱타임 + 호가.""" tp = str(tick_primary or live_tick_primary()).strip().lower() or "kiwoom" op = str(ob_primary or live_ob_primary()).strip().lower() or "kiwoom" alt_t = "kiwoom" if tp == "kis" else "kis" alt_o = "kiwoom" if op == "kis" else "kis" parts: List[str] = [ f"틱1차설정={tp}(체인 {tp}→{alt_t}→ls)", f"호가1차설정={op}(체인 {op}→{alt_o}→ls)", ] tr = tick_rec or {} if tr.get("label") or tr.get("vendor"): lab = str(tr.get("label") or format_vendor_label( str(tr.get("vendor") or ""), int(tr.get("tier") or tick_feed_tier(str(tr.get("vendor") or ""), tp)), bool(tr.get("spilled")), )) bit = [f"틱실제={lab}"] if tr.get("price"): try: bit.append("px=%s" % int(float(tr.get("price") or 0))) except (TypeError, ValueError): bit.append("px=%s" % tr.get("price")) if tr.get("tick_time"): bit.append("t=%s" % tr.get("tick_time")) if tr.get("lag_sec") is not None: bit.append("lag=%ss" % tr.get("lag_sec")) elif tr.get("age_ms") is not None: bit.append("age=%sms" % tr.get("age_ms")) parts.append(" ".join(bit)) else: parts.append("틱실제=없음") obr = ob_rec or {} if obr.get("label") or obr.get("vendor"): lab = str(obr.get("label") or format_vendor_label( str(obr.get("vendor") or ""), int(obr.get("tier") or tick_feed_tier(str(obr.get("vendor") or ""), op)), bool(obr.get("spilled")), )) bit = [f"호가실제={lab}"] if obr.get("best_bid") or obr.get("best_ask"): bit.append("bid=%s ask=%s" % (obr.get("best_bid") or 0, obr.get("best_ask") or 0)) if obr.get("spread_pct") is not None: bit.append("spr=%s%%" % obr.get("spread_pct")) if obr.get("bid_ask_ratio") is not None: bit.append("or=%s" % obr.get("bid_ask_ratio")) if obr.get("snap_time"): bit.append("snap=%s" % obr.get("snap_time")) if obr.get("ob_source"): bit.append("src=%s" % obr.get("ob_source")) if obr.get("age_ms") is not None: bit.append("age=%sms" % obr.get("age_ms")) parts.append(" ".join(bit)) else: parts.append("호가실제=없음") return " | ".join(parts) def _tick_second_key(tick: Dict[str, Any]) -> str: tt = str(tick.get("tick_time") or "")[:14] if len(tt) >= 14: return tt[:14] if len(tt) >= 12: return tt[:12] + "00" return tt def merge_ticks_time_axis_fallback( ticks: Sequence[Dict[str, Any]], *, main_src: str = "kiwoom", max_lag_sec: Optional[float] = None, ) -> List[Dict[str, Any]]: """같은 초에는 메인만. 메인·보조 모두 lag>나이면 그 초는 비움(OHLC 메우지 않음). 분봉에 메인 1건 있다고 보조를 통째로 버리지 않음 (실매 2초 실패→2차와 동일). """ if not ticks: return list(ticks) main = str(main_src or "kiwoom").strip().lower() if main not in ("kis", "kiwoom"): main = "kiwoom" age = live_feed_fallback_max_age_sec() if max_lag_sec is None else float(max_lag_sec) grouped: Dict[str, List[Dict[str, Any]]] = {} order: List[str] = [] for t in ticks: if not isinstance(t, dict): continue sk = _tick_second_key(t) if not sk: continue if sk not in grouped: grouped[sk] = [] order.append(sk) grouped[sk].append(t) out: List[Dict[str, Any]] = [] for sk in order: bucket = grouped[sk] mains: List[Dict[str, Any]] = [] aux: List[Dict[str, Any]] = [] for t in bucket: src = str(t.get("source") or "").strip().lower() lag = t.get("_lag_sec") if is_feed_read_stale(lag, age): continue if src == main: mains.append(t) else: aux.append(t) if mains: out.extend(mains) else: out.extend(aux) return out def merge_ls_ticks_third_fallback( ticks_by_code: Dict[str, Dict[str, List[Dict[str, Any]]]], ls_ticks_by_code: Dict[str, Dict[str, List[Dict[str, Any]]]], *, max_lag_sec: Optional[float] = None, ) -> int: """실매 틱 3차: 같은 초에 kis/kiwoom 이 이미 있으면 LS 안 넣음. 메인·2차가 비었거나(또는 나이컷으로 비움) 그 초만 ``ls_ws_ticks`` 로 채움. ``LIVE_FEED_FALLBACK_MAX_AGE_SEC``(기본 3) — lag 모르면 버리지 않음(실매와 동일). 반환=추가한 틱 수. """ if not ls_ticks_by_code: return 0 age = live_feed_fallback_max_age_sec() if max_lag_sec is None else float(max_lag_sec) added = 0 for code, ls_minutes in ls_ticks_by_code.items(): code_s = str(code or "").strip() if not code_s or not isinstance(ls_minutes, dict): continue main_minutes = ticks_by_code.setdefault(code_s, {}) occupied: set = set() for _mk, ticks in main_minutes.items(): for t in ticks or []: if not isinstance(t, dict): continue sk = _tick_second_key(t) if sk: occupied.add(sk) for mk, ls_ticks in ls_minutes.items(): mk_s = str(mk or "").strip() if not mk_s: continue bucket = main_minutes.setdefault(mk_s, []) for t in ls_ticks or []: if not isinstance(t, dict): continue sk = _tick_second_key(t) if not sk or sk in occupied: continue lag = t.get("_lag_sec") if is_feed_read_stale(lag, age): continue row = dict(t) row["source"] = "ls" bucket.append(row) occupied.add(sk) added += 1 if bucket: bucket.sort(key=lambda x: str(x.get("tick_time") or "")) return added def orderbook_row_lag_seconds(row: Dict[str, Any]) -> Optional[int]: """recv_ts vs snap_time 초 차이. 파싱 실패면 None(유효로 봄).""" recv_ts = str(row.get("recv_ts") or "").strip() st = str(row.get("snap_time") or "").strip() if len(recv_ts) < 19 or len(st) < 6: return None try: recv_dt = datetime.strptime(recv_ts[:19], "%Y-%m-%d %H:%M:%S") except Exception: return None try: if len(st) >= 14 and st[:14].isdigit(): pkt = datetime.strptime(st[:14], "%Y%m%d%H%M%S") elif len(st) >= 6 and st[-6:].isdigit(): pkt = datetime.strptime(recv_ts[:10].replace("-", "") + st[-6:], "%Y%m%d%H%M%S") else: return None return int((recv_dt - pkt).total_seconds()) except Exception: return None def candle_garbage_fallback_enabled() -> bool: """쓰레기 봉 스킵(2차·3차 소스 폴백) 활성 여부. 2026-09-06: 기본 True 복원. bar_is_garbage 를 wall-clock(recv_ts) 기준으로 정정 후 실매 RAM 3초컷과 동일 논리로 동작 → 실매↔백테 정합. (docs/정합성.md §9) BACKTEST_CANDLE_GARBAGE_OFF=1 은 실험용 강제 OFF (일반 운영은 사용 금지). """ import os if os.environ.get("BACKTEST_CANDLE_GARBAGE_OFF") == "1": return False try: return bool(get_env_bool("CANDLE_GARBAGE_FALLBACK", True)) except Exception: return True def bar_end_datetime(candle_time: str, tf_min: int): """봉 시작 candle_time(YYYYMMDDHHMM) + tf → 봉 끝 datetime.""" from kis_trader.engine.candle_rollup import add_candle_minutes end = add_candle_minutes(str(candle_time or "")[:12], int(tf_min or 1)) if not end or len(end) < 12: return None try: return datetime.strptime(end[:12], "%Y%m%d%H%M") except Exception: return None def tick_in_bar_bucket(tick_time: str, candle_time: str, tf_min: int) -> bool: from kis_trader.engine.candle_rollup import add_candle_minutes ct = str(candle_time or "").strip()[:12] raw = str(tick_time or "").strip().replace(":", "").replace("-", "").replace(" ", "") if len(ct) < 12: return False tmin = raw[:12] if len(raw) >= 12 else "" if len(tmin) < 12 or not tmin.isdigit(): return False end = add_candle_minutes(ct, int(tf_min or 1)) if not end: return False return ct <= tmin < end[:12] def bar_is_garbage( ticks: Sequence[Dict[str, Any]], *, candle_time: str, tf_min: int, source: str, missing_policy: str = "hole", ) -> bool: """그 분·그 source 틱이 없거나 전부 수신시각(recv_ts) 대비 읽기나이 초과면 True(구멍). 2026-09-06 정정 (docs/정합성.md): 기준 시각을 **봉끝** → **각 틱의 recv_ts** 로 변경. 이유: 실매 RAM 3초컷은 `wall-clock(수신순간) - tick_time` 기준. 이전 봉끝 기준은 유동성 낮은 종목·봉 시작 근처 틱이 부당 스킵되어 실매와 갭 발생. 이제 실매·백테 동일 wall-clock lag 판정 → 정합. fallback: recv_ts 없으면 옛 봉끝 기준(레거시 데이터 호환). missing_policy: hole — 0건이면 쓰레기 (옵투나 공책) keep — 0건이면 유지 (실매 링이 그 분을 커버 못 할 때 호출측에서 keep) 증권사 시계끼리 비교하지 않음. """ if not candle_garbage_fallback_enabled(): return False src = str(source or "").strip().lower() bucket: List[Dict[str, Any]] = [] for t in ticks or []: if not isinstance(t, dict): continue tsrc = str(t.get("source") or "").strip().lower() if tsrc and src and tsrc != src: continue raw = t.get("tick_time_raw") or t.get("tick_time") or "" if tick_in_bar_bucket(str(raw), candle_time, tf_min): bucket.append(t) if not bucket: return str(missing_policy or "hole").strip().lower() == "hole" bar_end = bar_end_datetime(candle_time, tf_min) for t in bucket: # 1) 사전 계산된 _lag_sec 가 있으면 최우선 (로더가 recv_ts - tick_time 을 이미 계산) # None 이면 판정 불가 → 이 틱은 안전(False), 다음 폴백으로. pre_lag = t.get("_lag_sec") if pre_lag is not None: try: _pl = float(pre_lag) if not is_feed_read_stale(_pl): return False # _pl 이 stale 이면 다음 틱으로 계속. 이 틱은 컷 대상. continue except (TypeError, ValueError): pass # 아래 폴백으로 raw = str(t.get("tick_time_raw") or t.get("tick_time") or "") # 2) 각 틱의 wall-clock 기준(recv_ts) 대비 lag = 실매 RAM 판정과 동일 논리 recv_dt: Optional[datetime] = None rts = t.get("recv_ts") if rts is not None and str(rts).strip(): try: if isinstance(rts, datetime): recv_dt = rts else: s = str(rts).strip().replace("T", " ") recv_dt = datetime.strptime(s[:19], "%Y-%m-%d %H:%M:%S") except Exception: recv_dt = None # 3) recv_ts 없거나 파싱 실패 → 봉끝 기준(레거시 폴백, 실매보다 보수적) base_dt = recv_dt if recv_dt is not None else bar_end if base_dt is None: # 봉끝도 파싱 실패면 판정 포기 → False(버리지 않음) return False lag = packet_lag_seconds(raw, now_dt=base_dt) if not is_feed_read_stale(lag): return False return True def live_bar_is_garbage( ticks: Sequence[Dict[str, Any]], *, candle_time: str, tf_min: int, source: str, ) -> bool: """실매 링버퍼: 그 소스 틱이 그 분까지 없으면 판정 불가 → False(유지). LS 틱은 TickRecorder 링에 안 들어갈 수 있음. 다른 증권사 틱만으로 cover 하면 LS/2차 봉을 구멍으로 잘못 버린다. 소스별 cover 만 본다. """ if not candle_garbage_fallback_enabled(): return False ct = str(candle_time or "")[:12] if len(ct) < 12: return False src = str(source or "").strip().lower() cover = False for t in ticks or []: if not isinstance(t, dict): continue tsrc = str(t.get("source") or "").strip().lower() if tsrc and src and tsrc != src: continue raw = str(t.get("tick_time_raw") or t.get("tick_time") or "").strip() digits = raw.replace(":", "").replace("-", "").replace(" ", "") tmin = digits[:12] if len(digits) >= 12 else "" if tmin and tmin.isdigit() and tmin <= ct: cover = True break if not cover: return False return bar_is_garbage( ticks, candle_time=candle_time, tf_min=tf_min, source=source, missing_policy="hole", )