diff --git a/docs/rust_engine_parity_port_plan.md b/docs/rust_engine_parity_port_plan.md index 130afdc..8180e70 100644 --- a/docs/rust_engine_parity_port_plan.md +++ b/docs/rust_engine_parity_port_plan.md @@ -400,6 +400,12 @@ python3 -u scripts/test_live_execution_validation.py | A3 | reaper 훅 — mode_refine 완료 시 phase1/2 JSON 자동 `register_result_json_as_job` (결과 잡아이디 UI 원복) | `kis_trader/backtest/optuna_web_jobs.py` | | A4 | 옵투나 잡 삭제 API/UI 신설 (`/api/optuna/delete/`, `/api/optuna/delete_bulk`, JS 🗑 버튼) | `optuna_web_jobs.py`, `backtest_web.py`, `static/js/backtest.js` | | A5 | whipsaw 표준 키 통일 (`whipsaw_enabled` 옛 alias 완전 제거 → `whipsaw_filter_enabled` 단일) | 6개 파일 (룰 28) | +| A6 | `bar_is_garbage` wall-clock(recv_ts) 정합 + `bt_candle_source` 로더 통일 (breakout_tick_loader 위임 · 22.8% → 5.2%) | `bt_candle_source.py`, `feed_fallback.py`, `breakout_tick_loader.py` | +| A7 | 파케이 초고속 로드 게이트 해제 (`BACKTEST_USE_RUST` 무관 항상 사용) — I/O 최적화이지 엔진 계산 아님 | `breakout_tick_loader.py:277-282` | +| A8 | 4전략 Optuna 파일에서 Rust 세션 초기화·`_rust_session_id` 세팅을 `BACKTEST_USE_RUST==1` 로 게이팅 (Python 라디오 = 진짜 Python) | `param_search_optuna.py`, `optuna_scalping.py`, `optuna_breakout.py`, `optuna_momentum.py` | +| A9 | `optuna_web_jobs.py` 의 `BACKTEST_USE_RUST` 세팅 반전 로직 정정 (`use_rust=True → "1"`, False → "0" 명시) | `optuna_web_jobs.py:3763` | +| A10 | `annotate_optuna_period_daily_avg` 가 결과 JSON 에 `use_rust` 자동 기록 → auto-register 뱃지 ❔ 재발 방지 | `optuna_common.py:112` | +| A11 | `KIWOOM_TICK_LIVE_MAX_LAG_SEC` 폐기 → 3벤더 공통 `WS_TICK_DB_SAVE_LAG_CUT_ENABLED` 스위치 (기본 OFF=전부 저장, 통계 유지) | `kiwoom_ws.py`, `kis_ws.py`, `ls_ws.py`, `database.py` | ## 부록 B: 발견된 Rust 정합 이슈 12건 (우선순위) @@ -420,4 +426,59 @@ python3 -u scripts/test_live_execution_validation.py --- -**끝.** 이 문서 기준으로 새 대화에서 P0 부터 착수. +## 부록 C: D안 — Rust 엔진 데이터 절반짜리 상태 (2026-09-06 발견 · P0 최우선) + +### C.1 사실 (실측) + +**Rust 세션 초기화** (`kis_rust_core.init_backtest_session_json(session_id, candles_json)`) 는 **candles(봉) 만** 넘긴다. +`run_engine_trial_scalp/breakout/momentum/tail(session, params)` 도 신규 데이터 없이 파라미터만 받는다. + +즉 **현재 Rust 엔진은 데이터 절반짜리**: + +| 데이터 | Python 엔진 (`engine_params`) | Rust 세션 | +|-------|:---:|:---:| +| 봉 (`codes_candles`) | ✅ | ✅ | +| 틱 (`ticks_by_code`) | ✅ | ❌ **미전달** | +| 호가 (`_bt_orderbook_by_code`) | ✅ | ❌ **미전달** | +| 프로그램매매 (`_bt_program_by_code`) | ✅ | ❌ **미전달** | + +- 확인 위치: `optuna_scalping.py:428-435`, `optuna_breakout.py:428-435`, `optuna_momentum.py:492-500`, `param_search_optuna.py:518-525` +- Rust 소비처: `scalping_backtest_common.py:337`, `breakout_backtest_common.py`, `momentum_backtest_common.py`, `tail_backtest_common.py` 모두 `if rust_session:` 분기에서 params 만 넘김 + +### C.2 의미 + +Python 엔진이 정합 개선(A6~A11) 을 받아도, **Rust 세션은 여전히 봉만 보고 판정** → 다음이 그대로 갭: + +1. **틱 청산** (WS 체결가·max_price 갱신·트레일 스톱) → Rust 는 종가 OHLC 기반 근사 +2. **호가 필터** (진입/exit/stop 3구간) → Rust 는 무필터 (post_filter 만 있음) +3. **프로그램매매 필터** → Rust 완전 없음 +4. **whipsaw 필터 (진입 시점)** → Rust 는 진입 이후 post_filter (실제 진입 자체를 못 막음) + +이 상태에서 Optuna Rust 결과를 실매에 적용하면 **실매는 위 4개 필터로 진입 자체가 안 되는 후보** 가 Rust 상 최적화 대상이 되어 랭킹 왜곡. + +### C.3 P0 로드맵 (D안 · 사용자 승인 시 착수) + +| 단계 | 파일 | 작업 | +|-----|-----|------| +| D-1 | `kis_rust_src/src/session_manager.rs` | `Session` 구조체에 `ticks: HashMap>>` (code→minute→ticks), `orderbook_snaps: HashMap>`, `program_ticks: HashMap>` 필드 추가 | +| D-2 | `kis_rust_src/src/lib.rs` | `init_backtest_session_json` 시그니처 확장: `(session_id, candles_json, ticks_json, orderbook_json, program_json)`. 옛 시그니처는 wrapper 로 유지 (하위호환) | +| D-3 | `kis_rust_src/src/{scalp,breakout,momentum,tail}.rs` | `try_sell_on_ticks` Rust 포팅 (Python `resolve_backtest_sell` 미러). 봉 high 로 max_price 선반영 금지 (룰 §부록 B #3) | +| D-4 | `kis_rust_src/src/filters.rs` (신설) | `orderbook_reject_for_entry`, `whipsaw_filter_hit`, `program_filter_reject` Rust 포팅. Python 과 1원 단위 일치 golden test | +| D-5 | 4전략 Optuna 파일 | `init_backtest_session_json` 호출 시 `ticks_by_code`, `orderbook_map`, `program_map` 모두 JSON 직렬화해서 전달 | +| D-6 | golden test (`scripts/golden_run.py`) | Rust ON + Python OFF 결과 diff = 0 (같은 필터·같은 데이터에서) | + +### C.4 리스크 + +- Rust 크레이트 대공사 (`session_manager.rs` 필드 대량 추가 · 새 소비 로직) +- `ticks_by_code` 전체 JSON 직렬화 → 200종목 1일치 ~100MB. **JSON 대신 shared memory 또는 Arrow IPC 고려 필요** (P0-6 서브태스크) +- 이번 세션에서 착수하지 않음. 이 문서는 **P0 최우선 항목 (부록 B #3 대체·확장)** 으로 다음 세션 시작점 + +### C.5 임시 대응 (D안 완료 전까지) + +- Optuna 실행은 **Python 라디오 기본** 유지 (A8 로 이미 가능) +- Rust 는 봉만 필요한 단순 백테·golden 진행에만 사용 +- Optuna best 를 실매 적용 전 반드시 **웹 백테(Python) 로 재검증** (룰 15) + +--- + +**끝.** 이 문서 기준으로 새 대화에서 **부록 C P0 (D안 · 데이터 완전 전달)** 부터 착수. diff --git a/kis_trader/backtest/bt_candle_source.py b/kis_trader/backtest/bt_candle_source.py index 54d8987..0330d9b 100644 --- a/kis_trader/backtest/bt_candle_source.py +++ b/kis_trader/backtest/bt_candle_source.py @@ -97,37 +97,50 @@ def _ws_ticks_table(market: Optional[str]) -> str: return "ws_ticks" -def _ticks_for_bar_garbage( - db, - code: str, - start_key: str, - end_key: str, - *, - market: Optional[str] = None, -) -> Optional[List[Dict[str, Any]]]: - from kis_trader.engine.feed_fallback import candle_garbage_fallback_enabled +def _ingest_flat_tick_rows( + rows: List[Dict[str, Any]], + target: Dict[str, List[Dict[str, Any]]], +) -> int: + """SELECT 행 → 종목별 flat 리스트 (봉 pick 쓰레기검사용). - if not candle_garbage_fallback_enabled(): - return None - sk, ek = _tick_time_bounds(start_key, end_key) - mk = (market or "").strip().upper() - table = _ws_ticks_table(mk if mk in ("US", "KR") else "KR") - try: - if mk in ("US", "KR"): - rows = db.conn.execute( - f"SELECT tick_time, source, tick_time_raw FROM {table} " - "WHERE market=%s AND code=%s AND tick_time >= %s AND tick_time <= %s", - (mk, str(code).strip(), sk, ek), - ).fetchall() - else: - rows = db.conn.execute( - f"SELECT tick_time, source, tick_time_raw FROM {table} " - "WHERE market=%s AND code=%s AND tick_time >= %s AND tick_time <= %s", - ("KR", str(code).strip(), sk, ek), - ).fetchall() - return [dict(r) for r in (rows or [])] - except Exception: - return None + ``breakout_tick_loader._ingest_tick_rows`` 는 minute_key 버킷 딕트를 만들지만, + bt_candle_source 소비자 (``dedupe_by_read_pairs``) 는 flat 리스트를 원한다. + → 같은 컬럼 (``recv_ts`` + ``_lag_sec``) 을 유지하되 리턴 형식만 다르게. + """ + from kis_trader.backtest.breakout_tick_loader import _tick_row_lag_seconds + + n = 0 + for r in rows or []: + code = str(r.get("code") or "").strip() + if not code: + continue + tt = str(r.get("tick_time") or "")[:14] + if len(tt) < 12: + continue + kw_lag = _tick_row_lag_seconds(r) + # recv_ts: bar_is_garbage 가 wall-clock 정합 판정에 사용 (docs/정합성.md §9) + # 문자열로 보존. None 이면 봉끝 폴백. + _recv_raw = r.get("recv_ts") + _recv_str = "" + if _recv_raw is not None: + try: + _recv_str = ( + _recv_raw.strftime("%Y-%m-%d %H:%M:%S") + if hasattr(_recv_raw, "strftime") + else str(_recv_raw) + ) + except Exception: + _recv_str = str(_recv_raw) + target[code].append({ + "code": code, + "tick_time": tt, + "source": r.get("source") or "", + "tick_time_raw": str(r.get("tick_time_raw") or ""), + "_lag_sec": kw_lag, + "recv_ts": _recv_str, + }) + n += 1 + return n def _load_ticks_by_code_bulk( @@ -138,49 +151,64 @@ def _load_ticks_by_code_bulk( market: Optional[str] = None, codes_filter: Optional[Sequence[str]] = None, ) -> Optional[Dict[str, List[Dict[str, Any]]]]: - """쓰레기 검사용 틱 — 기간 1~2쿼리. OFF 면 None (봉만 dedupe).""" + """쓰레기 검사용 틱 — 로더 통일 (2026-09-06 C안). + + 이전(자체 SELECT): ``tick_time, source, tick_time_raw`` 만 → recv_ts/_lag_sec 미포함 + → ``bar_is_garbage`` 가 봉끝 폴백 → **22.8% 대량 컷** (실매와 다른 기준). + + 현재: ``breakout_tick_loader._fetch_ws_ticks_day_rows`` 위임 → + 엔진 로더와 **동일 SELECT + 동일 lag 계산 + 동일 recv_ts** → + 쓰레기 판정도 wall-clock 정합. + + - KR → ``ws_ticks`` / US → ``ws_ticks_us`` + - OFF 면 None (봉만 dedupe) + - 리턴 형식은 flat ``Dict[code, List[tick]]`` (bt_candle_source 소비자 호환) + """ from kis_trader.engine.feed_fallback import candle_garbage_fallback_enabled if not candle_garbage_fallback_enabled(): return None - sk, ek = _tick_time_bounds(start_key, end_key) - mk = (market or "").strip().upper() - out: Dict[str, List[Dict[str, Any]]] = defaultdict(list) - want: List[str] = [] - if codes_filter: - want = [str(c).strip() for c in codes_filter if str(c).strip()] - in_sql = "" - in_params: List[str] = [] - if want: - in_sql = " AND code IN (" + ",".join(["%s"] * len(want)) + ")" - in_params = want - def _pull(table: str, mkt: str) -> int: - rows = db.conn.execute( - f"SELECT code, tick_time, source, tick_time_raw FROM {table} " - "WHERE market=%s AND tick_time >= %s AND tick_time <= %s" - + in_sql, - (mkt, sk, ek, *in_params), - ).fetchall() - n = 0 - for r in rows or []: - code = str(r.get("code") or "").strip() - if not code: + from kis_trader.backtest.breakout_tick_loader import ( + _candle_keys_to_tick_range, + _fetch_ws_ticks_day_rows, + _iter_tick_day_chunks, + _ws_ticks_table as _bl_ws_ticks_table, + ) + + mk = (market or "").strip().upper() + tt_start, tt_end = _candle_keys_to_tick_range(start_key, end_key) + + want: Optional[set] = None + if codes_filter: + want = {str(c).strip() for c in codes_filter if str(c).strip()} + + out: Dict[str, List[Dict[str, Any]]] = defaultdict(list) + + def _pull(table: str, mkt_str: str) -> int: + total = 0 + for chunk_s, chunk_e in _iter_tick_day_chunks(tt_start, tt_end): + try: + rows = _fetch_ws_ticks_day_rows(table, mkt_str, chunk_s, chunk_e, want) + except Exception as e: + logger.warning( + "쓰레기검사 틱 로드 실패 (table=%s day=%s): %s", + table, chunk_s[:8], e, + ) continue - out[code].append(dict(r)) - n += 1 - return n + total += _ingest_flat_tick_rows(rows, out) + return total try: n_kr = n_us = 0 if mk == "US": - n_us = _pull("ws_ticks_us", "US") + n_us = _pull(_bl_ws_ticks_table("US"), "US") elif mk == "KR": - n_kr = _pull("ws_ticks", "KR") + n_kr = _pull(_bl_ws_ticks_table("KR"), "KR") else: - n_kr = _pull("ws_ticks", "KR") + n_kr = _pull(_bl_ws_ticks_table("KR"), "KR") try: - n_us = _pull("ws_ticks_us", "US") + n_us = _pull(_bl_ws_ticks_table("US"), "US") except Exception: n_us = 0 logger.info("📥 틱 bulk 쓰레기검사용: KR=%s US=%s 종목=%s", n_kr, n_us, len(out)) @@ -190,6 +218,34 @@ def _load_ticks_by_code_bulk( return None +def _ticks_for_bar_garbage( + db, + code: str, + start_key: str, + end_key: str, + *, + market: Optional[str] = None, +) -> Optional[List[Dict[str, Any]]]: + """단일 종목 쓰레기 검사용 틱 — 로더 통일 (bulk 재사용, 2026-09-06 C안). + + 이전엔 별도 SELECT (``tick_time, source, tick_time_raw`` 만) → recv_ts/_lag_sec 없음. + 이제 ``_load_ticks_by_code_bulk`` (breakout_tick_loader 위임) 를 코드 1개로 호출 → + 동일 정합. + """ + from kis_trader.engine.feed_fallback import candle_garbage_fallback_enabled + + if not candle_garbage_fallback_enabled(): + return None + if not str(code or "").strip(): + return None + result = _load_ticks_by_code_bulk( + db, start_key, end_key, market=market, codes_filter=[str(code).strip()], + ) + if result is None: + return None + return result.get(str(code).strip(), []) + + def list_ws_candle_codes( db, timeframe: int,