feat: Enhance trading system with new permanent subscription features and order book management
Changes: - Added a new API endpoint for managing permanent subscriptions, allowing users to enable or disable subscriptions dynamically. - Implemented a function to fill candle data from Kiwoom, ensuring that only relevant data is inserted into the database. - Introduced a mechanism to handle master subscription states, improving the management of subscription statuses. - Updated the database schema to include new fields for managing subscription states and order book filtering. Impact: - These enhancements improve the flexibility and reliability of the trading system, allowing for better management of subscriptions and order book data, while reducing the risk of data inconsistencies. 히스토리 align 제거 븅신같은 초기설계 아예 제거 진입모드에 구멍메움 호가진입을 켜도 호가가 안들어올때 호가 안보고 그냥 사버림
This commit is contained in:
@@ -1,103 +1,78 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
백테·Optuna ws_candles 소스 선택 — 실매 LIVE_TICK_PROVIDER 와 동일 우선순위.
|
||||
백테·Optuna ws_candles 소스 선택 — 실매 get_candles 와 동일.
|
||||
|
||||
UI/CLI ``CANDLE_SOURCE``:
|
||||
- 빈값(「기본」): LIVE_TICK_PROVIDER 순서로 candle_time 디듑
|
||||
kiwoom 메인 → kiwoom, kis, rest, rollup_1m
|
||||
kis 메인 → kis, kiwoom, rest, rollup_1m
|
||||
- ``kis`` / ``kiwoom``: 해당 소스만 (디버그·비교용)
|
||||
읽기: 틱 메인 WS 한 칸, 없으면 kiwoom+rest 구멍, 그다음 kiwoom+rollup.
|
||||
CANDLE_SOURCE=kis|kiwoom 이면 그 메인(+구멍). 빈값이면 LIVE_TICK_PROVIDER.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from typing import Any, Dict, List, Optional, Sequence, Tuple
|
||||
|
||||
from kis_trader.utils.env import get_env_from_db
|
||||
from kis_trader.ws.candle_series import (
|
||||
dedupe_by_read_pairs,
|
||||
live_read_label,
|
||||
live_read_pairs,
|
||||
normalize_source_channel,
|
||||
)
|
||||
|
||||
# 실매 CandleAggregator._ALL_CANDLE_SOURCES 와 동일
|
||||
ReadPair = Tuple[str, str]
|
||||
|
||||
# 레거시 호환 이름 (실제 필터는 live_read_pairs)
|
||||
BT_WS_CANDLE_SOURCES: Tuple[str, ...] = ("kiwoom", "kis", "rest", "rollup_1m")
|
||||
|
||||
|
||||
def resolve_bt_candle_source_override() -> str:
|
||||
"""'' = 병합 모드, 'kis'|'kiwoom' = 단일 소스."""
|
||||
raw = os.environ.get("CANDLE_SOURCE")
|
||||
if raw is None:
|
||||
try:
|
||||
raw = get_env_from_db("CANDLE_SOURCE", "")
|
||||
except Exception:
|
||||
raw = ""
|
||||
s = (str(raw or "")).strip().lower()
|
||||
if s in ("kis", "kiwoom"):
|
||||
return s
|
||||
return ""
|
||||
"""'' = LIVE 병합, 'kis'|'kiwoom' = 그 메인(+키움 REST 구멍)."""
|
||||
from kis_trader.ws.candle_series import _provider_and_override
|
||||
_provider, override = _provider_and_override()
|
||||
return override
|
||||
|
||||
|
||||
def live_candle_source_order() -> Tuple[str, ...]:
|
||||
"""실매 get_candles 병합 순서 (LIVE_TICK_PROVIDER)."""
|
||||
try:
|
||||
provider = (
|
||||
get_env_from_db("LIVE_TICK_PROVIDER", "kiwoom") or "kiwoom"
|
||||
).strip().lower()
|
||||
except Exception:
|
||||
provider = "kiwoom"
|
||||
if provider == "kis":
|
||||
return ("kis", "kiwoom", "rest", "rollup_1m")
|
||||
return ("kiwoom", "kis", "rest", "rollup_1m")
|
||||
"""호환: source 문자열만. 실제 조회는 resolve_bt_read_pairs."""
|
||||
return tuple(dict.fromkeys(p[0] for p in live_read_pairs()))
|
||||
|
||||
|
||||
def resolve_bt_read_pairs() -> Tuple[ReadPair, ...]:
|
||||
return live_read_pairs()
|
||||
|
||||
|
||||
def resolve_bt_candle_source_order() -> Tuple[str, ...]:
|
||||
"""백테/Optuna 조회에 쓸 소스 순서."""
|
||||
override = resolve_bt_candle_source_override()
|
||||
if override:
|
||||
return (override,)
|
||||
"""호환 래퍼 — 신규 코드는 resolve_bt_read_pairs 사용."""
|
||||
return live_candle_source_order()
|
||||
|
||||
|
||||
def resolve_bt_candle_source_label() -> str:
|
||||
"""상태 표시용 — 'KIS'|'KIWOOM'|'LIVE(kiwoom)' 등."""
|
||||
override = resolve_bt_candle_source_override()
|
||||
if override:
|
||||
return override.upper()
|
||||
try:
|
||||
provider = (
|
||||
get_env_from_db("LIVE_TICK_PROVIDER", "kiwoom") or "kiwoom"
|
||||
).strip().lower()
|
||||
except Exception:
|
||||
provider = "kiwoom"
|
||||
return f"LIVE({provider})"
|
||||
return live_read_label()
|
||||
|
||||
|
||||
def dedupe_candle_rows(
|
||||
rows: Sequence[Dict[str, Any]],
|
||||
source_order: Optional[Sequence[str]] = None,
|
||||
source_order: Optional[Sequence[Any]] = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""candle_time 기준 디듑 — 앞 소스 우선. ``source`` 컬럼은 결과에서 제거."""
|
||||
order = tuple(source_order or resolve_bt_candle_source_order())
|
||||
rank = {s: i for i, s in enumerate(order)}
|
||||
seen: Dict[str, Tuple[int, Dict[str, Any]]] = {}
|
||||
for row in rows:
|
||||
ct = str(row.get("candle_time") or "")[:12]
|
||||
if not ct:
|
||||
continue
|
||||
src = str(row.get("source") or "kis").strip().lower()
|
||||
pri = rank.get(src, 999)
|
||||
prev = seen.get(ct)
|
||||
if prev is None or pri < prev[0]:
|
||||
seen[ct] = (pri, dict(row))
|
||||
out: List[Dict[str, Any]] = []
|
||||
for ct in sorted(seen.keys()):
|
||||
d = seen[ct][1]
|
||||
d.pop("source", None)
|
||||
out.append(d)
|
||||
return out
|
||||
"""candle_time 1행. source_order 가 쌍이면 그대로, 아니면 live_read_pairs."""
|
||||
pairs: Optional[Sequence[ReadPair]] = None
|
||||
if source_order:
|
||||
first = source_order[0]
|
||||
if isinstance(first, (tuple, list)) and len(first) >= 2:
|
||||
pairs = tuple((str(a), str(b)) for a, b in source_order) # type: ignore[misc]
|
||||
else:
|
||||
# 레거시 source 문자열 → 정규화 후 순위
|
||||
mapped: List[ReadPair] = []
|
||||
for s in source_order:
|
||||
mapped.append(normalize_source_channel(str(s), ""))
|
||||
pairs = tuple(mapped)
|
||||
return dedupe_by_read_pairs(rows, pairs)
|
||||
|
||||
|
||||
def _source_in_sql(sources: Sequence[str]) -> Tuple[str, List[str]]:
|
||||
if len(sources) == 1:
|
||||
return " AND source=%s", [sources[0]]
|
||||
ph = ",".join(["%s"] * len(sources))
|
||||
return f" AND source IN ({ph})", list(sources)
|
||||
def _pairs_in_sql(pairs: Sequence[ReadPair]) -> Tuple[str, List[str]]:
|
||||
parts: List[str] = []
|
||||
params: List[str] = []
|
||||
for src, ch in pairs:
|
||||
parts.append("(source=%s AND channel=%s)")
|
||||
params.extend([src, ch])
|
||||
return " AND (" + " OR ".join(parts) + ")", params
|
||||
|
||||
|
||||
def list_ws_candle_codes(
|
||||
@@ -108,9 +83,9 @@ def list_ws_candle_codes(
|
||||
*,
|
||||
market: Optional[str] = None,
|
||||
) -> List[str]:
|
||||
"""기간 내 종목 코드 — 선택된 소스 기준 DISTINCT."""
|
||||
sources = resolve_bt_candle_source_order()
|
||||
src_sql, src_params = _source_in_sql(sources)
|
||||
"""기간 내 종목 코드 — 선택된 (source,channel) 기준 DISTINCT."""
|
||||
pairs = resolve_bt_read_pairs()
|
||||
src_sql, src_params = _pairs_in_sql(pairs)
|
||||
mk = (market or "").strip().upper()
|
||||
if mk:
|
||||
rows = db.conn.execute(
|
||||
@@ -143,36 +118,15 @@ def fetch_ws_candles_for_code(
|
||||
market: Optional[str] = None,
|
||||
confirmed_only: bool = True,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""단일 종목·기간 봉 로드 — 소스 필터/병합 적용."""
|
||||
sources = resolve_bt_candle_source_order()
|
||||
"""단일 종목·기간 봉 로드 — 메인 WS 우선, 구멍 kiwoom+rest(+rollup)."""
|
||||
pairs = resolve_bt_read_pairs()
|
||||
confirmed_sql = " AND is_confirmed=1" if confirmed_only else ""
|
||||
mk = (market or "").strip().upper()
|
||||
|
||||
if len(sources) == 1:
|
||||
cols = f"candle_time, open, high, low, close, volume{peak_sel}{extra_select}"
|
||||
src = sources[0]
|
||||
if mk:
|
||||
rows = db.conn.execute(
|
||||
f"SELECT {cols} FROM ws_candles "
|
||||
"WHERE timeframe=%s AND code=%s AND market=%s "
|
||||
"AND candle_time >= %s AND candle_time <= %s"
|
||||
+ confirmed_sql
|
||||
+ " AND source=%s ORDER BY candle_time ASC",
|
||||
[int(timeframe), code, mk, start_key, end_key, src],
|
||||
).fetchall()
|
||||
else:
|
||||
rows = db.conn.execute(
|
||||
f"SELECT {cols} FROM ws_candles "
|
||||
"WHERE timeframe=%s AND code=%s "
|
||||
"AND candle_time >= %s AND candle_time <= %s"
|
||||
+ confirmed_sql
|
||||
+ " AND source=%s ORDER BY candle_time ASC",
|
||||
[int(timeframe), code, start_key, end_key, src],
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
cols = f"candle_time, open, high, low, close, volume, source{peak_sel}{extra_select}"
|
||||
src_sql, src_params = _source_in_sql(sources)
|
||||
src_sql, src_params = _pairs_in_sql(pairs)
|
||||
cols = (
|
||||
f"candle_time, open, high, low, close, volume, source, channel"
|
||||
f"{peak_sel}{extra_select}"
|
||||
)
|
||||
if mk:
|
||||
rows = db.conn.execute(
|
||||
f"SELECT {cols} FROM ws_candles "
|
||||
@@ -193,7 +147,7 @@ def fetch_ws_candles_for_code(
|
||||
+ " ORDER BY candle_time ASC",
|
||||
[int(timeframe), code, start_key, end_key, *src_params],
|
||||
).fetchall()
|
||||
return dedupe_candle_rows([dict(r) for r in rows], sources)
|
||||
return dedupe_by_read_pairs([dict(r) for r in rows], pairs)
|
||||
|
||||
|
||||
def fetch_ws_candles_warmup_before(
|
||||
@@ -210,24 +164,14 @@ def fetch_ws_candles_warmup_before(
|
||||
"""기간 시작 이전 N봉 — prepend 웜업용 (오래된→최신)."""
|
||||
if limit <= 0 or not before_candle_time:
|
||||
return []
|
||||
sources = resolve_bt_candle_source_order()
|
||||
pairs = resolve_bt_read_pairs()
|
||||
confirmed_sql = " AND is_confirmed=1" if confirmed_only else ""
|
||||
fetch_limit = max(int(limit) * max(len(sources), 1), int(limit) + 50)
|
||||
|
||||
if len(sources) == 1:
|
||||
cols = f"candle_time, open, high, low, close, volume{peak_sel}{extra_select}"
|
||||
rows = db.conn.execute(
|
||||
f"SELECT {cols} FROM ws_candles "
|
||||
"WHERE timeframe=%s AND code=%s AND candle_time < %s"
|
||||
+ confirmed_sql
|
||||
+ " AND source=%s ORDER BY candle_time DESC LIMIT %s",
|
||||
[int(timeframe), code, before_candle_time, sources[0], fetch_limit],
|
||||
).fetchall()
|
||||
bars = [dict(r) for r in reversed(rows)]
|
||||
return bars[-limit:] if len(bars) > limit else bars
|
||||
|
||||
cols = f"candle_time, open, high, low, close, volume, source{peak_sel}{extra_select}"
|
||||
src_sql, src_params = _source_in_sql(sources)
|
||||
fetch_limit = max(int(limit) * max(len(pairs), 1), int(limit) + 50)
|
||||
src_sql, src_params = _pairs_in_sql(pairs)
|
||||
cols = (
|
||||
f"candle_time, open, high, low, close, volume, source, channel"
|
||||
f"{peak_sel}{extra_select}"
|
||||
)
|
||||
rows = db.conn.execute(
|
||||
f"SELECT {cols} FROM ws_candles "
|
||||
"WHERE timeframe=%s AND code=%s AND candle_time < %s"
|
||||
@@ -236,5 +180,5 @@ def fetch_ws_candles_warmup_before(
|
||||
+ " ORDER BY candle_time DESC LIMIT %s",
|
||||
[int(timeframe), code, before_candle_time, *src_params, fetch_limit],
|
||||
).fetchall()
|
||||
bars = dedupe_candle_rows([dict(r) for r in rows], sources)
|
||||
bars = dedupe_by_read_pairs([dict(r) for r in reversed(rows)], pairs)
|
||||
return bars[-limit:] if len(bars) > limit else bars
|
||||
|
||||
Reference in New Issue
Block a user