Files
kis_bot/kis_trader/execution/account_cash.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

175 lines
6.5 KiB
Python

"""
kis_trader/execution/account_cash.py — 가용 예수금 메모리 캐시 + kv_store 영속
================================================================================
- 매수체크 루프마다 REST 호출하지 않음.
- main._fetch_asset_snapshot() / 장 시작·heartbeat 에서 API 동기화.
- 체결 후 메모리에서 증감, kv_store 에 주기적·체결 시 저장 (재시작 복구).
kv_store 키 (잡다한 런타임 상태 — 별도 테이블 불필요):
account.available_cash : 주문가능 예수금(원, 정수 문자열)
account.cash_updated_at : ISO 시각
account.cash_source : api | trade_delta | startup
"""
from __future__ import annotations
import threading
import time
from datetime import datetime as dt
from typing import Any, Dict, Optional
from ..utils.env import get_env_float, get_env_int
from ..utils.logger import LOG_YELLOW, LOG_RESET, get_logger
logger = get_logger("kis_trader.account_cash")
KV_AVAILABLE_CASH = "account.available_cash"
KV_CASH_UPDATED_AT = "account.cash_updated_at"
KV_CASH_SOURCE = "account.cash_source"
class AccountCashLedger:
"""프로세스 내 가용 예수금 — OrderManager 매수 직전 부족 시에만 qty 조정에 사용."""
def __init__(self, db) -> None:
# TradeDBExt — get_kv/set_kv 는 raw TradeDB
self._db = db
self._raw = getattr(db, "raw", db)
self._lock = threading.Lock()
self._cash: float = 0.0
self._has_baseline: bool = False
self._last_persist_ts: float = 0.0
self._load_from_kv()
def _load_from_kv(self) -> None:
try:
v = self._raw.get_kv(KV_AVAILABLE_CASH)
if v is None or str(v).strip() == "":
return
cash = float(str(v).replace(",", "").strip())
if cash > 0:
with self._lock:
self._cash = cash
self._has_baseline = True
logger.info(
"📂 [예수금캐시] kv_store 복원: %s원 (updated=%s)",
f"{cash:,.0f}",
self._raw.get_kv(KV_CASH_UPDATED_AT) or "?",
)
except Exception as e:
logger.debug("예수금 kv 복원 실패: %s", e)
def available_cash(self) -> float:
with self._lock:
return float(self._cash)
def has_baseline(self) -> bool:
with self._lock:
return self._has_baseline and self._cash > 0
def sync_from_snapshot(self, snap: Dict[str, Any], source: str = "api") -> None:
"""KIS inquire-balance 파싱 결과(cash=D+2 or dnca) → 메모리·kv."""
try:
cash = float(snap.get("cash", 0) or 0)
except (TypeError, ValueError):
return
if cash <= 0:
return
with self._lock:
self._cash = cash
self._has_baseline = True
self._persist(source=source, force=False)
def invalidate_baseline(self, reason: str = "startup_sync_failed") -> None:
"""
stale kv baseline 사용 중단.
- 재시작 직후 API 동기화 실패 시 이전 계좌/구좌의 캐시 오염을 막기 위해 사용.
- _cash 값은 남겨두되 baseline 플래그를 내려 clamp 로직이 개입하지 않게 한다.
"""
with self._lock:
self._has_baseline = False
try:
self._raw.set_kv(KV_CASH_SOURCE, reason)
except Exception:
pass
def apply_trade_delta(self, delta_krw: float, source: str = "trade_delta") -> None:
"""체결 후 낙관적 증감 (매수 음수, 매도 양수). baseline 없으면 스킵."""
if abs(delta_krw) < 1:
return
with self._lock:
if not self._has_baseline:
return
self._cash = max(0.0, self._cash + delta_krw)
self._persist(source=source, force=True)
def _persist(self, source: str, force: bool) -> None:
interval = max(5, get_env_int("ACCOUNT_CASH_PERSIST_SEC", 60))
now = time.time()
if not force and (now - self._last_persist_ts) < interval:
return
with self._lock:
cash = self._cash
try:
self._raw.set_kv(KV_AVAILABLE_CASH, str(int(round(cash))))
self._raw.set_kv(KV_CASH_UPDATED_AT, dt.now().strftime("%Y-%m-%d %H:%M:%S"))
self._raw.set_kv(KV_CASH_SOURCE, source)
self._last_persist_ts = now
except Exception as e:
logger.debug("예수금 kv 저장 실패: %s", e)
def clamp_buy_qty_if_insufficient(
self,
qty: int,
price_ref: float,
code: str,
name: str,
strategy_id: str,
) -> tuple[int, Optional[str]]:
"""
기존 전략이 계산한 qty 를 우선.
주문금액(수수료 버퍼 포함) > 가용예수금 일 때만 예수금 % 로 상한 재계산.
Returns:
(adjusted_qty, reason) — qty=0 이면 매수 스킵, reason=insufficient_cash
"""
if qty <= 0 or price_ref <= 0:
return qty, None
if not self.has_baseline():
return qty, None
fee_buf = max(1.0, get_env_float("ORDER_CASH_FEE_BUFFER", 1.01))
order_amt = qty * price_ref * fee_buf
available = self.available_cash()
if order_amt <= available:
return qty, None
pct = max(0.01, min(1.0, get_env_float("ORDER_CASH_PCT", 0.95)))
cap_cash = available * pct
divide = get_env_int("ORDER_CASH_DIVIDE_BY_MAX_STOCKS", 1)
if divide > 0:
max_stocks = max(1, get_env_int("MAX_STOCKS", 4))
cap_cash = cap_cash / max_stocks
unit_cost = price_ref * fee_buf
max_qty = int(cap_cash / unit_cost) if unit_cost > 0 else 0
if max_qty < 1:
logger.warning(
"%s🚫 [예수금부족] [%s] %s %s: 필요 %s원 > 가용 %s원 → 1주도 불가%s",
LOG_YELLOW, strategy_id, name, code,
f"{order_amt:,.0f}", f"{available:,.0f}", LOG_RESET,
)
return 0, "insufficient_cash"
if max_qty < qty:
logger.warning(
"%s💰 [예수금조정] [%s] %s %s: qty %d%d "
"(주문 %s원 > 가용 %s원, cap=%s원×%.0f%%)%s",
LOG_YELLOW, strategy_id, name, code, qty, max_qty,
f"{order_amt:,.0f}", f"{available:,.0f}",
f"{cap_cash:,.0f}", pct * 100, LOG_RESET,
)
return max_qty, "cash_clamped"
return qty, None