feat: Add DART strategy and related configurations
ㅇ Changes: - Introduced the DART strategy to the trading system, including its configuration and integration into the existing framework. - Updated the database schema to include DART-specific tables for disclosures and watchlists. - Enhanced the backtesting and parameter search functionalities to support the DART strategy. - Implemented new rules for browser verification and API interactions to ensure compliance with the updated DART strategy. Impact: - These additions expand the trading capabilities of the system, allowing for more comprehensive analysis and execution of DART-related strategies, while maintaining system integrity and performance.
This commit is contained in:
@@ -1124,6 +1124,91 @@ class CandleAggregator:
|
||||
self._flush_batch(batch)
|
||||
logger.info("CandleDBWriter 스레드 종료")
|
||||
|
||||
@staticmethod
|
||||
def _ws_candle_freeze_on_confirm() -> bool:
|
||||
"""확정봉(is_confirmed=1) OHLCV 동결 — docs/정합성.md. 기본 true."""
|
||||
return get_env_bool("WS_CANDLE_FREEZE_ON_CONFIRM", True)
|
||||
|
||||
def _load_confirmed_ohlcv_from_db(
|
||||
self, code: str, tf: int, candle_times: list,
|
||||
) -> Dict[str, Dict]:
|
||||
"""
|
||||
freeze 재시작 정합: RAM 이 비어도 DB 에 이미 확정된 봉은 REST 로 덮지 않고
|
||||
DB 값을 RAM 에 시드한다. (실매 RAM ≠ DB 재발 방지)
|
||||
"""
|
||||
out: Dict[str, Dict] = {}
|
||||
if not self.db or not candle_times:
|
||||
return out
|
||||
times = sorted({str(t)[:12] for t in candle_times if str(t)[:12]})
|
||||
if not times:
|
||||
return out
|
||||
# IN 절 길이 제한 — 청크
|
||||
chunk_n = max(50, int(get_env_int("WS_CANDLE_FREEZE_DB_LOOKUP_CHUNK", 200)))
|
||||
try:
|
||||
for i in range(0, len(times), chunk_n):
|
||||
chunk = times[i : i + chunk_n]
|
||||
ph = ",".join(["%s"] * len(chunk))
|
||||
rows = self.db.conn.execute(
|
||||
f"""
|
||||
SELECT candle_time, `open`, high, low, close, volume,
|
||||
rsi_2, rsi_3, rsi_5, source, holding_peak
|
||||
FROM ws_candles
|
||||
WHERE code=%s AND timeframe=%s AND is_confirmed=1
|
||||
AND candle_time IN ({ph})
|
||||
""",
|
||||
(code, int(tf), *chunk),
|
||||
).fetchall()
|
||||
for r in rows or []:
|
||||
ct = str(r.get("candle_time") or "")[:12]
|
||||
if not ct:
|
||||
continue
|
||||
out[ct] = dict(r)
|
||||
except Exception as e:
|
||||
logger.debug("freeze DB lookup 실패(%s %dM): %s", code, tf, e)
|
||||
return out
|
||||
|
||||
@staticmethod
|
||||
def _ws_candles_upsert_sql(*, freeze: bool) -> str:
|
||||
"""
|
||||
freeze ON: 이미 확정된 행의 OHLCV/RSI/volume/source 유지.
|
||||
미확정(is_confirmed=0)→확정·갱신은 허용. holding_peak 만 항상 GREATEST.
|
||||
freeze OFF: 기존처럼 덮어쓰기.
|
||||
"""
|
||||
if freeze:
|
||||
dup = """
|
||||
ON DUPLICATE KEY UPDATE
|
||||
`open`=IF(is_confirmed=1, `open`, VALUES(`open`)),
|
||||
high=IF(is_confirmed=1, high, GREATEST(high, VALUES(high))),
|
||||
low=IF(is_confirmed=1, low, VALUES(low)),
|
||||
close=IF(is_confirmed=1, close, VALUES(close)),
|
||||
volume=IF(is_confirmed=1, volume, VALUES(volume)),
|
||||
rsi_2=IF(is_confirmed=1, rsi_2, VALUES(rsi_2)),
|
||||
rsi_3=IF(is_confirmed=1, rsi_3, VALUES(rsi_3)),
|
||||
rsi_5=IF(is_confirmed=1, rsi_5, VALUES(rsi_5)),
|
||||
is_confirmed=IF(is_confirmed=1, 1, VALUES(is_confirmed)),
|
||||
source=IF(is_confirmed=1, source, VALUES(source)),
|
||||
holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)),
|
||||
updated_at=IF(is_confirmed=1, updated_at, VALUES(updated_at))
|
||||
"""
|
||||
else:
|
||||
dup = """
|
||||
ON DUPLICATE KEY UPDATE
|
||||
`open`=VALUES(`open`), high=GREATEST(high, VALUES(high)),
|
||||
low=VALUES(low), close=VALUES(close), volume=VALUES(volume),
|
||||
rsi_2=VALUES(rsi_2), rsi_3=VALUES(rsi_3), rsi_5=VALUES(rsi_5),
|
||||
is_confirmed=VALUES(is_confirmed),
|
||||
holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)),
|
||||
updated_at=VALUES(updated_at)
|
||||
"""
|
||||
return f"""
|
||||
INSERT INTO ws_candles
|
||||
(code, timeframe, candle_time, `open`, high, low, close,
|
||||
volume, rsi_2, rsi_3, rsi_5, is_confirmed, source, holding_peak, updated_at)
|
||||
VALUES
|
||||
(%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
{dup}
|
||||
"""
|
||||
|
||||
def _flush_batch(self, batch: list) -> None:
|
||||
"""
|
||||
배치 리스트를 DB에 한 번의 executemany 로 INSERT.
|
||||
@@ -1154,22 +1239,12 @@ class CandleAggregator:
|
||||
item.get("holding_peak"),
|
||||
now_str,
|
||||
))
|
||||
sql = self._ws_candles_upsert_sql(
|
||||
freeze=self._ws_candle_freeze_on_confirm(),
|
||||
)
|
||||
# ws_candles 테이블이 존재하면 배치 INSERT (없으면 조용히 skip)
|
||||
self.db.conn.execute(
|
||||
"""
|
||||
INSERT INTO ws_candles
|
||||
(code, timeframe, candle_time, `open`, high, low, close,
|
||||
volume, rsi_2, rsi_3, rsi_5, is_confirmed, source, holding_peak, updated_at)
|
||||
VALUES
|
||||
(%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
`open`=VALUES(`open`), high=GREATEST(high, VALUES(high)),
|
||||
low=VALUES(low), close=VALUES(close), volume=VALUES(volume),
|
||||
rsi_2=VALUES(rsi_2), rsi_3=VALUES(rsi_3), rsi_5=VALUES(rsi_5),
|
||||
is_confirmed=VALUES(is_confirmed),
|
||||
holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)),
|
||||
updated_at=VALUES(updated_at)
|
||||
""",
|
||||
sql,
|
||||
rows[0],
|
||||
) if len(rows) == 1 else self._executemany_batch(rows)
|
||||
logger.debug("💾 [배치저장] %d봉 → DB", len(rows))
|
||||
@@ -1178,22 +1253,10 @@ class CandleAggregator:
|
||||
|
||||
def _executemany_batch(self, rows: list) -> None:
|
||||
"""여러 봉을 executemany 로 한 번에 INSERT."""
|
||||
sql = """
|
||||
INSERT INTO ws_candles
|
||||
(code, timeframe, candle_time, `open`, high, low, close,
|
||||
volume, rsi_2, rsi_3, rsi_5, is_confirmed, source, holding_peak, updated_at)
|
||||
VALUES
|
||||
(%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
`open`=VALUES(`open`), high=GREATEST(high, VALUES(high)),
|
||||
low=VALUES(low), close=VALUES(close), volume=VALUES(volume),
|
||||
rsi_2=VALUES(rsi_2), rsi_3=VALUES(rsi_3), rsi_5=VALUES(rsi_5),
|
||||
is_confirmed=VALUES(is_confirmed),
|
||||
holding_peak=GREATEST(COALESCE(holding_peak,0), COALESCE(VALUES(holding_peak),0)),
|
||||
updated_at=VALUES(updated_at)
|
||||
"""
|
||||
sql = self._ws_candles_upsert_sql(
|
||||
freeze=self._ws_candle_freeze_on_confirm(),
|
||||
)
|
||||
# pymysql executemany: cursor.executemany(sql, list_of_tuples)
|
||||
import pymysql
|
||||
with self.db.conn._lock:
|
||||
self.db.conn._ensure_connected()
|
||||
cur = self.db.conn._conn.cursor()
|
||||
@@ -1230,6 +1293,16 @@ class CandleAggregator:
|
||||
import datetime as _dt
|
||||
return _dt.datetime.now().strftime("%Y%m%d%H%M")
|
||||
|
||||
@staticmethod
|
||||
def _open_bucket_ctime(tf: int, now=None) -> str:
|
||||
"""현재 시각 기준 진행 중(미완성) 봉의 candle_time (YYYYMMDDHHMM)."""
|
||||
import datetime as _dt
|
||||
from kis_trader.engine.candle_rollup import floor_candle_time_to_tf
|
||||
|
||||
dt0 = now if now is not None else _dt.datetime.now()
|
||||
raw = dt0.strftime("%Y%m%d%H%M")
|
||||
return floor_candle_time_to_tf(raw, int(tf) or 1)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# RSI 계산 (확정 봉 close 리스트 기반)
|
||||
# ------------------------------------------------------------------
|
||||
@@ -1344,8 +1417,9 @@ class CandleAggregator:
|
||||
[트랙 1] RSI 계산 + ``_confirmed`` RAM 버퍼 적재 (매수 루프 즉시 참조용)
|
||||
[트랙 2] DB 기록 Queue 적재 (논블로킹)
|
||||
|
||||
동일 ``candle_time`` 이 이미 있으면 **append 하지 않고 upsert**.
|
||||
volume 은 더 큰 쪽을 유지 (갭보정 불완전봉 < WS 누적 < REST 완전봉).
|
||||
동일 ``candle_time`` 이 이미 있으면:
|
||||
- ``WS_CANDLE_FREEZE_ON_CONFIRM``(기본 true): OHLCV 유지(첫 확정 승), holding_peak 만 갱신
|
||||
- freeze OFF: volume 더 큰 쪽 upsert (레거시)
|
||||
"""
|
||||
code, tf = key
|
||||
ctime = str(cur.get("candle_time") or "")[:12]
|
||||
@@ -1362,20 +1436,24 @@ class CandleAggregator:
|
||||
)
|
||||
|
||||
hp = self._holding_peak_for(code)
|
||||
freeze = self._ws_candle_freeze_on_confirm()
|
||||
if idx >= 0:
|
||||
old = buf[idx]
|
||||
old_vol = int(old.get("volume") or 0)
|
||||
# 이미 더 완전한 volume 이 있으면(예: REST 완전봉) WS 부분봉으로 덮지 않음
|
||||
if new_vol < old_vol:
|
||||
# freeze: 이미 확정된 봉은 OHLCV 고정 (REST/재확정이 키우지 않음)
|
||||
if freeze or new_vol < old_vol:
|
||||
confirmed_candle = dict(old)
|
||||
confirmed_candle["is_confirmed"] = 1
|
||||
if hp is not None:
|
||||
confirmed_candle["holding_peak"] = max(
|
||||
float(confirmed_candle.get("holding_peak") or 0), float(hp),
|
||||
)
|
||||
confirmed_candle["high"] = max(
|
||||
float(confirmed_candle.get("high") or 0), float(hp),
|
||||
)
|
||||
# freeze 시 high 도 동결 — peak 만 메타로 보관
|
||||
if not freeze:
|
||||
confirmed_candle["high"] = max(
|
||||
float(confirmed_candle.get("high") or 0), float(hp),
|
||||
)
|
||||
buf[idx] = confirmed_candle
|
||||
return confirmed_candle
|
||||
|
||||
low_cands = [
|
||||
@@ -1512,11 +1590,19 @@ class CandleAggregator:
|
||||
if rest_df is None or rest_df.empty:
|
||||
return 0
|
||||
|
||||
# 진행 중 분봉은 confirmed 에 넣지 않음 — merge_confirmed_bars 에서도 재필터.
|
||||
skip_incomplete = get_env_bool("WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET", True)
|
||||
open_bucket = self._open_bucket_ctime(tf) if skip_incomplete else ""
|
||||
|
||||
rows: list = []
|
||||
skipped_open = 0
|
||||
for _, row in rest_df.iterrows():
|
||||
ctime = str(row.get("time", ""))[:12]
|
||||
if not ctime or len(ctime) < 12:
|
||||
continue
|
||||
if open_bucket and ctime >= open_bucket:
|
||||
skipped_open += 1
|
||||
continue
|
||||
close = float(row.get("close", 0) or 0)
|
||||
if close <= 0:
|
||||
continue
|
||||
@@ -1529,6 +1615,11 @@ class CandleAggregator:
|
||||
"volume": int(float(row.get("volume", 0) or 0)),
|
||||
"source": "rest",
|
||||
})
|
||||
if skipped_open:
|
||||
logger.info(
|
||||
"⏭ [갭보정] %s %dM 진행분(>=%s) %d봉 confirmed 제외 (매수 직전봉 왜곡 방지)",
|
||||
code, tf, open_bucket, skipped_open,
|
||||
)
|
||||
return self.merge_confirmed_bars(code, tf, rows, log_tag="REST")
|
||||
|
||||
def merge_confirmed_bars(
|
||||
@@ -1538,20 +1629,68 @@ class CandleAggregator:
|
||||
bars: list,
|
||||
*,
|
||||
log_tag: str = "merge",
|
||||
skip_incomplete_bucket: Optional[bool] = None,
|
||||
now=None,
|
||||
) -> int:
|
||||
"""
|
||||
확정봉 리스트를 RAM(+DB 큐)에 병합.
|
||||
|
||||
- 신규 candle_time → insert
|
||||
- 기존 candle_time → volume 이 더 클 때만 OHLCV upsert
|
||||
(갭보정 불완전봉을 REST/완전 롤업이 덮어쓰도록)
|
||||
- 기존 candle_time → ``WS_CANDLE_FREEZE_ON_CONFIRM``(기본 true) 이면 skip
|
||||
(freeze OFF: volume 더 클 때만 OHLCV upsert — 레거시)
|
||||
- ``WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET``(기본 true): 진행 중 버킷
|
||||
(``candle_time >= 현재 봉시작``) 은 insert/update 하지 않고,
|
||||
이미 RAM 에 있으면 제거. 장초 REST 미완성봉 → 직전봉% 몸통 왜곡 방지.
|
||||
"""
|
||||
if not bars:
|
||||
return 0
|
||||
# bars 비어도 진행분 purge 는 수행 (재갭보정 정리)
|
||||
pass
|
||||
|
||||
if skip_incomplete_bucket is None:
|
||||
skip_incomplete_bucket = get_env_bool(
|
||||
"WS_GAP_FILL_SKIP_INCOMPLETE_BUCKET", True,
|
||||
)
|
||||
freeze = self._ws_candle_freeze_on_confirm()
|
||||
open_bucket = (
|
||||
self._open_bucket_ctime(tf, now=now) if skip_incomplete_bucket else ""
|
||||
)
|
||||
|
||||
# lock 밖에서 DB 조회 (재시작 후 RAM 공백 → REST 가 DB 확정봉을 덮는 것 방지)
|
||||
db_frozen: Dict[str, Dict] = {}
|
||||
if freeze and bars and self.db is not None:
|
||||
want_times = []
|
||||
for row in bars:
|
||||
ct = str(row.get("candle_time") or row.get("time") or "")[:12]
|
||||
if ct and len(ct) >= 12:
|
||||
want_times.append(ct)
|
||||
db_frozen = self._load_confirmed_ohlcv_from_db(code, tf, want_times)
|
||||
|
||||
with self._lock:
|
||||
key = (code, tf)
|
||||
closes = self._closes.setdefault(key, [])
|
||||
conf_buf = self._confirmed.setdefault(key, [])
|
||||
|
||||
# 이미 들어온 진행분 confirmed 제거 (갭보정 직후·롤업 공통)
|
||||
purged = 0
|
||||
if open_bucket and conf_buf:
|
||||
kept = []
|
||||
for c in conf_buf:
|
||||
ct = str(c.get("candle_time") or "")[:12]
|
||||
if ct and ct >= open_bucket:
|
||||
purged += 1
|
||||
continue
|
||||
kept.append(c)
|
||||
if purged:
|
||||
conf_buf[:] = kept
|
||||
closes[:] = [float(c.get("close") or 0) for c in conf_buf]
|
||||
logger.info(
|
||||
"🧹 [갭보정] %s %dM 진행분 confirmed %d봉 제거 (>=%s)",
|
||||
code, tf, purged, open_bucket,
|
||||
)
|
||||
|
||||
if not bars:
|
||||
return 0
|
||||
|
||||
by_time = {
|
||||
str(c.get("candle_time", ""))[:12]: i
|
||||
for i, c in enumerate(conf_buf)
|
||||
@@ -1559,10 +1698,16 @@ class CandleAggregator:
|
||||
}
|
||||
inserted = 0
|
||||
updated = 0
|
||||
skipped_freeze = 0
|
||||
seeded_db = 0
|
||||
skipped_open = 0
|
||||
for row in bars:
|
||||
ctime = str(row.get("candle_time") or row.get("time") or "")[:12]
|
||||
if not ctime or len(ctime) < 12:
|
||||
continue
|
||||
if open_bucket and ctime >= open_bucket:
|
||||
skipped_open += 1
|
||||
continue
|
||||
close = float(row.get("close", 0) or 0)
|
||||
if close <= 0:
|
||||
continue
|
||||
@@ -1571,6 +1716,10 @@ class CandleAggregator:
|
||||
idx = by_time.get(ctime)
|
||||
|
||||
if idx is not None:
|
||||
# freeze-on-confirm: 첫 확정본 유지 (갭보정/백필이 lookback 시리즈를 키우지 않음)
|
||||
if freeze:
|
||||
skipped_freeze += 1
|
||||
continue
|
||||
old = conf_buf[idx]
|
||||
old_vol = int(old.get("volume") or 0)
|
||||
if new_vol <= old_vol:
|
||||
@@ -1614,6 +1763,53 @@ class CandleAggregator:
|
||||
updated += 1
|
||||
continue
|
||||
|
||||
# RAM 에 없고 DB 에 확정봉이 있으면 → DB 값으로만 시드 (REST로 덮지 않음)
|
||||
if freeze and ctime in db_frozen:
|
||||
dr = db_frozen[ctime]
|
||||
d_close = float(dr.get("close") or 0)
|
||||
if d_close <= 0:
|
||||
d_close = close
|
||||
closes.append(d_close)
|
||||
if len(closes) > self._ram_buffer_max:
|
||||
closes.pop(0)
|
||||
rsi2 = dr.get("rsi_2")
|
||||
rsi3 = dr.get("rsi_3")
|
||||
rsi5 = dr.get("rsi_5")
|
||||
if rsi2 is None and rsi3 is None and rsi5 is None:
|
||||
rsi2, rsi3, rsi5 = self._compute_rsi_set(closes)
|
||||
candle = {
|
||||
"code": code,
|
||||
"tf": tf,
|
||||
"candle_time": ctime,
|
||||
"open": float(dr.get("open") or d_close),
|
||||
"high": float(dr.get("high") or d_close),
|
||||
"low": float(dr.get("low") or d_close),
|
||||
"close": d_close,
|
||||
"volume": int(float(dr.get("volume") or 0)),
|
||||
"rsi_2": rsi2,
|
||||
"rsi_3": rsi3,
|
||||
"rsi_5": rsi5,
|
||||
"is_confirmed": 1,
|
||||
"source": str(dr.get("source") or "db")[:10],
|
||||
}
|
||||
if dr.get("holding_peak") is not None:
|
||||
candle["holding_peak"] = dr.get("holding_peak")
|
||||
conf_buf.append(candle)
|
||||
by_time[ctime] = len(conf_buf) - 1
|
||||
if len(conf_buf) > self._ram_buffer_max:
|
||||
conf_buf.sort(key=lambda x: str(x.get("candle_time", "")))
|
||||
while len(conf_buf) > self._ram_buffer_max:
|
||||
conf_buf.pop(0)
|
||||
closes[:] = [float(c["close"]) for c in conf_buf]
|
||||
by_time = {
|
||||
str(c.get("candle_time", ""))[:12]: i
|
||||
for i, c in enumerate(conf_buf)
|
||||
if str(c.get("candle_time", ""))[:12]
|
||||
}
|
||||
# DB 이미 있음 → 쓰기 큐 불필요
|
||||
seeded_db += 1
|
||||
continue
|
||||
|
||||
closes.append(close)
|
||||
if len(closes) > self._ram_buffer_max:
|
||||
closes.pop(0)
|
||||
@@ -1650,16 +1846,21 @@ class CandleAggregator:
|
||||
except queue.Full:
|
||||
pass
|
||||
inserted += 1
|
||||
if inserted or updated:
|
||||
if inserted or updated or seeded_db:
|
||||
conf_buf.sort(key=lambda x: str(x.get("candle_time", "")))
|
||||
closes[:] = [float(c["close"]) for c in conf_buf]
|
||||
|
||||
if inserted or updated:
|
||||
logger.info(
|
||||
"🔧 [갭보정] %s %dM → %s insert=%d update=%d RAM+DB큐",
|
||||
code, tf, log_tag, inserted, updated,
|
||||
if skipped_open and log_tag:
|
||||
logger.debug(
|
||||
"⏭ [갭보정] %s %dM %s 진행분 %d봉 skip (>=%s)",
|
||||
code, tf, log_tag, skipped_open, open_bucket,
|
||||
)
|
||||
return inserted + updated
|
||||
if inserted or updated or skipped_freeze or seeded_db:
|
||||
logger.info(
|
||||
"🔧 [갭보정] %s %dM → %s insert=%d update=%d freeze_skip=%d db_seed=%d RAM+DB큐",
|
||||
code, tf, log_tag, inserted, updated, skipped_freeze, seeded_db,
|
||||
)
|
||||
return inserted + updated + seeded_db
|
||||
|
||||
def rollup_tf_from_1m(self, code: str, target_tf: int = 3) -> int:
|
||||
"""
|
||||
|
||||
Reference in New Issue
Block a user