Files
kis_bot/kis_trader/execution/orphan_reconcile.py
Hwang 61c72a8a4c feat(tests): 신규 키움 웹소켓 조건검색 및 실시간 조건검색 테스트 추가
변경 사항
----
- _test_kiwoom_condition_list.py: 키움 웹소켓 조건검색 '목록조회' 기능을 단독으로 테스트하는 스크립트 추가
- _test_kiwoom_condition_realtime.py: 'momentum' 조건식을 실시간으로 등록하고 초기 매칭 종목 리스트 및 실시간 편입/이탈을 수신하는 테스트 스크립트 추가
- _verify_columnar_bitid.py, _verify_shared_e2e_breakout.py, _verify_shared_e2e.py: 공유 메모리 및 dict 간의 데이터 일관성을 검증하는 테스트 추가

영향
----
- 신규 테스트 스크립트 추가로 키움 웹소켓 API의 기능 검증 및 안정성을 높임
- 기존 기능에 대한 영향 없음

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-06 01:27:00 +09:00

228 lines
7.8 KiB
Python

"""
kis_trader/execution/orphan_reconcile.py — 봇 고아 포지션 복구 (장마감 후 배치)
================================================================================
실계좌 잔고(inquire-balance, 페이징)와 active_trades·orders 를 대조해
「봇이 매수했는데 active_trades 미기록」 고아만 복구한다.
- 수동매수(orders BUY 없음) · MANUAL_HOLD_CODES 지정분은 건드리지 않음.
- 매매 알림은 보내지 않음(로그 + 선택적 요약 MM).
"""
from __future__ import annotations
import json
import logging
from typing import Any, Dict, List, Optional, Set, Tuple
from ..utils.env import get_env_bool, get_env_from_db
from ..utils.strategy_ids import canonical_strategy_id
logger = logging.getLogger("kis_trader.orphan_reconcile")
def parse_manual_hold_codes() -> Set[str]:
"""MANUAL_HOLD_CODES env — 쉼표/세미콜론 구분 종목코드."""
manual_raw = str(get_env_from_db("MANUAL_HOLD_CODES", "") or "")
return {
c.strip() for c in manual_raw.replace(";", ",").split(",") if c.strip()
}
def get_bot_bought_codes(db) -> Set[str]:
"""orders 테이블에 BUY 기록이 있는 종목코드."""
codes: Set[str] = set()
try:
cur = db.conn.execute("SELECT DISTINCT code FROM orders WHERE side='BUY'")
for row in (cur.fetchall() or []):
c = str(row.get("code") or "").strip()
if c:
codes.add(c)
except Exception as e:
logger.warning("orders BUY 코드 조회 실패: %s", e)
return codes
def get_portfolio_origin_sets(db) -> Tuple[Set[str], Set[str]]:
"""(bot_bought_codes, manual_hold_codes) — backtest_web 보유·매도 탭과 동일."""
return get_bot_bought_codes(db), parse_manual_hold_codes()
def _tracked_codes(db) -> Set[str]:
"""active_trades 에 한 건이라도 있는 종목코드 (전 전략)."""
out: Set[str] = set()
try:
cur = db.conn.execute("SELECT DISTINCT code FROM active_trades")
for row in (cur.fetchall() or []):
c = str(row.get("code") or "").strip()
if c:
out.add(c)
except Exception as e:
logger.warning("active_trades 코드 조회 실패: %s", e)
return out
def _latest_buy_order(db, code: str) -> Optional[Dict[str, Any]]:
"""해당 종목 최신 BUY orders 행 (복구 메타·체결가 참고)."""
try:
cur = db.conn.execute(
"""
SELECT * FROM orders
WHERE side='BUY' AND code=%s
ORDER BY submitted_at DESC
LIMIT 1
""",
(code,),
)
row = cur.fetchone()
return dict(row) if row else None
except Exception as e:
logger.warning("orders 최신 BUY 조회 실패 %s: %s", code, e)
return None
def _parse_order_meta(row: Dict[str, Any]) -> Dict[str, Any]:
"""orders.raw_json 에서 stop/target 등 복구."""
meta: Dict[str, Any] = {}
raw = row.get("raw_json")
if not raw:
return meta
try:
d = json.loads(raw) if isinstance(raw, str) else raw
if isinstance(d, dict):
meta = d
except Exception:
pass
return meta
def _num(v, default: float = 0.0) -> float:
try:
return float(str(v or 0).replace(",", "").strip())
except Exception:
return default
def reconcile_orphan_positions(order_mgr) -> Dict[str, Any]:
"""
장마감 후 1회 호출 — 봇 고아만 active_trades 에 upsert.
Returns:
reconciled, skipped_tracked, skipped_manual, skipped_no_order, failed
"""
result: Dict[str, Any] = {
"reconciled": [],
"skipped_tracked": [],
"skipped_manual": [],
"skipped_no_order": [],
"failed": [],
}
if not get_env_bool("ORPHAN_RECONCILE_ENABLED", True):
result["msg"] = "ORPHAN_RECONCILE_ENABLED=false"
return result
db = order_mgr.db
broker = order_mgr.get_broker_holdings(force=True)
if not getattr(order_mgr, "_holdings_last_fetch_ok", False):
result["error"] = "실계좌 잔고 조회 실패"
logger.warning("🧩 [고아복구] 잔고 API 실패 — 스킵")
return result
tracked = _tracked_codes(db)
bot_bought, manual_hold = get_portfolio_origin_sets(db)
for code, br in (broker or {}).items():
c = str(code).strip()
qty = int((br or {}).get("qty") or 0)
if qty <= 0:
continue
name = str((br or {}).get("name") or c)
if c in manual_hold:
result["skipped_manual"].append({"code": c, "name": name, "qty": qty})
continue
if c in tracked:
result["skipped_tracked"].append({"code": c, "name": name, "qty": qty})
continue
if c not in bot_bought:
result["skipped_no_order"].append({"code": c, "name": name, "qty": qty})
continue
order_row = _latest_buy_order(db, c)
if not order_row:
result["skipped_no_order"].append({"code": c, "name": name, "qty": qty})
continue
meta = _parse_order_meta(order_row)
strategy = canonical_strategy_id(
str(order_row.get("strategy_id") or meta.get("strategy_id") or "SCALP")
)
filled_px = _num(order_row.get("filled_avg_price"))
if filled_px <= 0:
filled_px = _num(order_row.get("price") or meta.get("price_ref"))
br_avg = _num((br or {}).get("avg_price"))
avg_buy = br_avg if br_avg > 0 else filled_px
cur_px = _num((br or {}).get("current_price")) or avg_buy
buy_date = (
str(order_row.get("filled_at") or order_row.get("submitted_at") or "").strip()
or None
)
stop_px = _num(meta.get("stop_price"))
target_px = _num(meta.get("target_price"))
atr_entry = _num(meta.get("atr_entry"))
trade_data = {
"code": c,
"name": str(order_row.get("name") or meta.get("name") or name),
"strategy": strategy,
"avg_buy_price": avg_buy,
"current_price": cur_px,
"stop_price": stop_px,
"target_price": target_px,
"max_price": max(avg_buy, cur_px),
"atr_entry": atr_entry,
"target_qty": qty,
"current_qty": qty,
"total_invested": avg_buy * qty,
"status": "HOLDING",
"buy_date": buy_date,
"size_class": meta.get("size_class") or "",
"entry_features": meta.get("entry_features") or {},
}
try:
ok = db.upsert_trade(trade_data)
if ok:
result["reconciled"].append({
"code": c,
"name": trade_data["name"],
"strategy": strategy,
"qty": qty,
"avg_buy": avg_buy,
})
logger.info(
"🧩 [고아복구] %s(%s) [%s] %d주 @ %.0f → active_trades",
trade_data["name"], c, strategy, qty, avg_buy,
)
else:
result["failed"].append({"code": c, "name": name, "error": "upsert_trade 실패"})
except Exception as e:
result["failed"].append({"code": c, "name": name, "error": str(e)})
logger.exception("고아복구 upsert 실패 %s", c)
try:
order_mgr.invalidate_holdings_cache()
except Exception:
pass
result["reconciled_count"] = len(result["reconciled"])
result["failed_count"] = len(result["failed"])
logger.info(
"🧩 [고아복구] 완료 — 복구 %d / 실패 %d / 보호(수동) %d / 이미추적 %d / 주문없음 %d",
len(result["reconciled"]),
len(result["failed"]),
len(result["skipped_manual"]),
len(result["skipped_tracked"]),
len(result["skipped_no_order"]),
)
return result