Files
kis_bot/kis_trader/execution/order_worker.py
Your Name 9ba9ab73b6 feat(backtest): 대대적인 Optuna 백테스트 웹 UI 및 백엔드 파이프라인 개편
- Web UI:
  - Optuna 탭 추가 및 mode_combo (최빈값 조합), 사후합격 Top 10 시각화 기능
  - 파라미터 분포(p25~p75, median, mode) 히스토그램 및 과적합(Overfit) 위험도 진단 UI 신설
  - 체크박스 렌더링 깨짐 현상을 네이티브(appearance: auto)로 강제 복구 (CSS)
  - 다단 트레일링 스탑, 꼬리 진입/돌파 손절 등 고급 조건 설정 폼 UI 고도화

- Backend (Optuna Jobs):
  - CLI 환경에서 구동된 Optuna json 결과물을 웹 대시보드로 읽어오는 import 기능 강화
  - JSON 메타데이터에 sort_by, mode, 호가 적용 여부 등 핵심 파라미터 파싱 누락 수정
  - optuna_mode_combo.py 등 최빈값 조합 및 후보군 2차 검증을 위한 신규 모듈 추가

- DB & Execution:
  - WebSocket 호가/틱 피드 수집 통계(api_feed_collect_stats) 메모리 캐시 최적화
  - KIS client 접속 키(approval_key) 등 인프라스트럭처 안정성 및 공유 관리 구조 개선
  - 테스트 및 디버깅용 briefing 마크다운 자동 생성 기능 추가
2026-09-01 02:47:51 +09:00

199 lines
6.6 KiB
Python

"""계좌 단일 주문 실행 Worker — 4전략 BUY/SELL intent 를 우선순위 FIFO 로 직렬 place."""
from __future__ import annotations
import itertools
import queue
import threading
import time
from typing import Any, Dict, Optional, TYPE_CHECKING
from ..utils.env import get_env_bool, get_env_int
from ..utils.logger import get_logger
from .order_intent import OrderIntent, resolve_order_priority
if TYPE_CHECKING:
from ..strategies.base import BaseStrategy
from .order_manager import OrderManager
logger = get_logger("kis_trader.order_worker")
_SENTINEL = object()
class AccountOrderWorker:
"""증권 계좌 1개에 맞춘 주문 대기줄 — 한 번에 place() 1건만."""
def __init__(self, order_mgr: "OrderManager") -> None:
self.order_mgr = order_mgr
self._pq: "queue.PriorityQueue[tuple]" = queue.PriorityQueue()
self._seq = itertools.count()
self._strategies: Dict[str, Any] = {}
self._running = False
self._thread: Optional[threading.Thread] = None
self._start_lock = threading.Lock()
self.enqueued_total = 0
self.processed_total = 0
self.dropped_total = 0
def register_strategy(self, strategy: "BaseStrategy") -> None:
sid = str(getattr(strategy, "strategy_id", "") or "").strip()
if sid:
self._strategies[sid] = strategy
def start(self) -> None:
with self._start_lock:
if self._thread is not None and self._thread.is_alive():
return
self._running = True
th = threading.Thread(
target=self._worker_loop,
name="AccountOrderWorker",
daemon=True,
)
self._thread = th
th.start()
logger.info("▶ [OrderWorker] 계좌 단일 주문 큐 기동")
def stop(self) -> None:
self._running = False
try:
self._pq.put_nowait((9999, next(self._seq), _SENTINEL))
except Exception:
pass
th = self._thread
if th is not None:
th.join(timeout=5.0)
def queue_depth(self) -> int:
try:
return int(self._pq.qsize())
except Exception:
return 0
def enqueue(
self,
strategy: "BaseStrategy",
side: str,
signal: Dict[str, Any],
*,
source: str = "scan",
priority: Optional[int] = None,
) -> bool:
"""intent 큐 적재. False=큐 만료 drop (inflight 는 호출자가 해제)."""
if not self._running:
return False
max_q = max(1, int(get_env_int("ORDER_WORKER_MAX_QUEUE", 500) or 500))
if self.queue_depth() >= max_q:
self.dropped_total += 1
code = str((signal or {}).get("code") or "")
logger.warning(
"⚠️ [OrderWorker] 큐 만료 drop [%s] %s %s depth>=%d",
getattr(strategy, "strategy_id", "?"),
side,
code,
max_q,
)
return False
prio = (
int(priority)
if priority is not None
else resolve_order_priority(side, signal or {}, source)
)
intent = OrderIntent(
priority=prio,
seq=next(self._seq),
strategy_id=str(getattr(strategy, "strategy_id", "") or ""),
side=str(side or "").upper(),
signal=dict(signal or {}),
source=str(source or "scan"),
)
self._pq.put((intent.priority, intent.seq, intent))
self.enqueued_total += 1
return True
def _worker_loop(self) -> None:
while True:
try:
_prio, _seq, item = self._pq.get(timeout=0.3)
except queue.Empty:
if not self._running:
break
continue
if item is _SENTINEL:
break
if not isinstance(item, OrderIntent):
continue
if not self._running:
self._release_inflight_for_intent(item)
continue
self._process_intent(item)
def _process_intent(self, intent: OrderIntent) -> None:
strat = self._strategies.get(intent.strategy_id)
code = str((intent.signal or {}).get("code") or "").strip()
side = str(intent.side or "").upper()
t0 = time.perf_counter()
try:
if strat is None:
logger.warning(
"⚠️ [OrderWorker] 미등록 전략 %s — skip %s %s",
intent.strategy_id,
side,
code,
)
return
if side == "SELL" and get_env_bool("REAL_BALANCE_VERIFY_BEFORE_SELL", True):
try:
self.order_mgr.prefetch_broker_holdings(force=False)
except Exception:
pass
if side == "BUY":
strat._submit_buy(intent.signal)
elif side == "SELL":
strat._submit_sell(intent.signal)
else:
logger.warning(
"⚠️ [OrderWorker] invalid side=%s [%s] %s",
side,
intent.strategy_id,
code,
)
except Exception as ex:
logger.warning(
"⚠️ [OrderWorker] place 예외 [%s] %s %s: %s",
intent.strategy_id,
side,
code or "-",
ex,
)
finally:
self._release_inflight_for_intent(intent)
self.processed_total += 1
elapsed_ms = (time.perf_counter() - t0) * 1000.0
if elapsed_ms >= 3000.0:
logger.warning(
"⚠️ [OrderWorker] 처리 지연 [%s] %s %s src=%s wait+place=%.0fms",
intent.strategy_id,
side,
code or "-",
intent.source,
elapsed_ms,
)
def _release_inflight_for_intent(self, intent: OrderIntent) -> None:
strat = self._strategies.get(intent.strategy_id)
if strat is None:
return
code = str((intent.signal or {}).get("code") or "").strip()
if not code:
return
side = str(intent.side or "").upper()
if side == "SELL":
release = getattr(strat, "_release_sell_inflight", None)
if callable(release):
release(code)
elif side == "BUY":
release = getattr(strat, "_release_buy_inflight", None)
if callable(release):
release(code)