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 제거 븅신같은 초기설계 아예 제거 진입모드에 구멍메움 호가진입을 켜도 호가가 안들어올때 호가 안보고 그냥 사버림
348 lines
12 KiB
Python
348 lines
12 KiB
Python
"""조건검색 진입/재진입/이탈 이벤트 → DB.
|
|
|
|
스냅샷 테이블(target/ls_candidates_history)은 **종목 목록**만 두고,
|
|
job(N/R/O·I/D) 은 이 테이블에 1이벤트 1행으로 적재한다.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from ..utils.env import get_env_bool
|
|
from ..utils.logger import get_logger
|
|
from .condition_common import JOB_KO, normalize_job_flag
|
|
|
|
logger = get_logger("kis_trader.cond_job")
|
|
|
|
CONDITION_JOB_EVENTS_DDL = """
|
|
CREATE TABLE IF NOT EXISTS condition_job_events (
|
|
id BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY,
|
|
event_time VARCHAR(30) NOT NULL,
|
|
strategy_id VARCHAR(32) NOT NULL DEFAULT '',
|
|
code VARCHAR(20) NOT NULL,
|
|
name VARCHAR(100) NOT NULL DEFAULT '',
|
|
job_flag VARCHAR(8) NOT NULL,
|
|
job_ko VARCHAR(16) NOT NULL DEFAULT '',
|
|
broker VARCHAR(16) NOT NULL,
|
|
source VARCHAR(32) NOT NULL DEFAULT '',
|
|
price DOUBLE NOT NULL DEFAULT 0,
|
|
broker_time VARCHAR(32) NOT NULL DEFAULT '',
|
|
query_name VARCHAR(64) NOT NULL DEFAULT '',
|
|
INDEX idx_cje_sid_et (strategy_id, event_time),
|
|
INDEX idx_cje_code_et (code, event_time),
|
|
INDEX idx_cje_et (event_time),
|
|
INDEX idx_cje_broker_job (broker, job_flag)
|
|
) CHARACTER SET utf8mb4
|
|
"""
|
|
|
|
# JOB_KO 는 condition_common 단일 소스 (여기서 재정의 금지)
|
|
|
|
|
|
def ensure_condition_job_events_table(db) -> bool:
|
|
if db is None:
|
|
return False
|
|
try:
|
|
db.conn.execute(CONDITION_JOB_EVENTS_DDL)
|
|
return True
|
|
except Exception as e:
|
|
logger.debug("condition_job_events DDL: %s", e)
|
|
return False
|
|
|
|
|
|
def insert_condition_job_event(
|
|
db,
|
|
*,
|
|
strategy_id: str,
|
|
code: str,
|
|
job_flag: str,
|
|
broker: str,
|
|
name: str = "",
|
|
price: float = 0.0,
|
|
broker_time: str = "",
|
|
query_name: str = "",
|
|
source: str = "",
|
|
event_time: Optional[str] = None,
|
|
) -> bool:
|
|
"""1건 INSERT. CONDITION_JOB_EVENTS_SAVE=false 면 스킵."""
|
|
if db is None:
|
|
return False
|
|
if not get_env_bool("CONDITION_JOB_EVENTS_SAVE", True):
|
|
return False
|
|
c = str(code or "").strip()
|
|
j = normalize_job_flag(job_flag)
|
|
if not c or not j:
|
|
return False
|
|
if j not in JOB_KO:
|
|
return False
|
|
sid = str(strategy_id or "").strip().upper()
|
|
br = str(broker or "").strip().lower() or "unknown"
|
|
et = event_time or datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
try:
|
|
px = float(price or 0)
|
|
except (TypeError, ValueError):
|
|
px = 0.0
|
|
ensure_condition_job_events_table(db)
|
|
try:
|
|
with db.conn:
|
|
db.conn.execute(
|
|
"""
|
|
INSERT INTO condition_job_events
|
|
(event_time, strategy_id, code, name, job_flag, job_ko,
|
|
broker, source, price, broker_time, query_name)
|
|
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
|
|
""",
|
|
(
|
|
et,
|
|
sid,
|
|
c,
|
|
str(name or c)[:100],
|
|
j,
|
|
JOB_KO.get(j, j),
|
|
br[:16],
|
|
str(source or "")[:32],
|
|
px,
|
|
str(broker_time or "")[:32],
|
|
str(query_name or "")[:64],
|
|
),
|
|
)
|
|
return True
|
|
except Exception as e:
|
|
logger.debug("condition_job_events INSERT 실패: %s", e)
|
|
return False
|
|
|
|
|
|
def query_condition_job_events(
|
|
db,
|
|
*,
|
|
start: str,
|
|
end: str,
|
|
strategy_id: str = "",
|
|
broker: str = "",
|
|
job_flag: str = "",
|
|
code: str = "",
|
|
limit: int = 500,
|
|
) -> List[Dict[str, Any]]:
|
|
"""웹/분석용 SELECT. start/end 는 YYYY-MM-DD 또는 datetime 문자열."""
|
|
if db is None:
|
|
return []
|
|
ensure_condition_job_events_table(db)
|
|
start_s = str(start or "").strip()
|
|
end_s = str(end or "").strip()
|
|
if len(start_s) == 10:
|
|
start_s = start_s + " 00:00:00"
|
|
if len(end_s) == 10:
|
|
end_s = end_s + " 23:59:59"
|
|
lim = max(1, min(int(limit or 500), 5000))
|
|
where = ["event_time >= %s", "event_time <= %s"]
|
|
params: List[Any] = [start_s, end_s]
|
|
sid = str(strategy_id or "").strip().upper()
|
|
if sid and sid != "ALL":
|
|
where.append("strategy_id = %s")
|
|
params.append(sid)
|
|
br = str(broker or "").strip().lower()
|
|
if br and br != "all":
|
|
where.append("broker = %s")
|
|
params.append(br)
|
|
jf = str(job_flag or "").strip().upper()
|
|
if jf and jf != "ALL":
|
|
where.append("job_flag = %s")
|
|
params.append(jf)
|
|
cd = str(code or "").strip()
|
|
if cd:
|
|
where.append("code LIKE %s")
|
|
params.append(cd + "%")
|
|
params.append(lim)
|
|
sql = (
|
|
"SELECT id, event_time, strategy_id, code, name, job_flag, job_ko, "
|
|
"broker, source, price, broker_time, query_name "
|
|
"FROM condition_job_events WHERE "
|
|
+ " AND ".join(where)
|
|
+ " ORDER BY event_time DESC, id DESC LIMIT %s"
|
|
)
|
|
try:
|
|
rows = db.conn.execute(sql, tuple(params)).fetchall()
|
|
except Exception as e:
|
|
logger.warning("condition_job_events SELECT: %s", e)
|
|
return []
|
|
out: List[Dict[str, Any]] = []
|
|
for r in rows or []:
|
|
if hasattr(r, "keys"):
|
|
out.append(dict(r))
|
|
else:
|
|
out.append({
|
|
"id": r[0], "event_time": r[1], "strategy_id": r[2],
|
|
"code": r[3], "name": r[4], "job_flag": r[5], "job_ko": r[6],
|
|
"broker": r[7], "source": r[8], "price": r[9],
|
|
"broker_time": r[10], "query_name": r[11],
|
|
})
|
|
return out
|
|
|
|
|
|
def build_condition_parity_report(db, *, day: str) -> Dict[str, Any]:
|
|
"""전략별 live소스·history·job 요약 (웹 탭용).
|
|
|
|
봇 RAM 은 웹에서 직접 못 읽으므로 **당일 history 최신 스냅샷** + job 집계로
|
|
구멍(LS history 공백·키움만 살아있음)을 가독성 있게 표시한다.
|
|
"""
|
|
from kis_trader.utils.env import get_env_from_db
|
|
|
|
day_s = str(day or "").strip()[:10]
|
|
if len(day_s) != 10:
|
|
day_s = datetime.now().strftime("%Y-%m-%d")
|
|
start_s = day_s + " 00:00:00"
|
|
end_s = day_s + " 23:59:59.999999"
|
|
sids = ("SCALP", "SHORT", "BREAKOUT", "MOMENTUM")
|
|
ensure_condition_job_events_table(db)
|
|
|
|
def _latest_hist(table: str, sid: str) -> Dict[str, Any]:
|
|
try:
|
|
row = db.conn.execute(
|
|
f"""
|
|
SELECT event_time AS et, COUNT(*) AS n
|
|
FROM {table}
|
|
WHERE strategy_id=%s AND event_time >= %s AND event_time <= %s
|
|
GROUP BY event_time
|
|
ORDER BY event_time DESC
|
|
LIMIT 1
|
|
""",
|
|
(sid, start_s, end_s),
|
|
).fetchone()
|
|
except Exception:
|
|
return {"event_time": None, "n": 0, "age_sec": None}
|
|
if not row:
|
|
return {"event_time": None, "n": 0, "age_sec": None}
|
|
et = row["et"] if hasattr(row, "keys") else row[0]
|
|
n = int(row["n"] if hasattr(row, "keys") else row[1] or 0)
|
|
age = None
|
|
try:
|
|
et_dt = datetime.strptime(str(et)[:19], "%Y-%m-%d %H:%M:%S")
|
|
age = max(0, int((datetime.now() - et_dt).total_seconds()))
|
|
except Exception:
|
|
age = None
|
|
return {"event_time": str(et) if et else None, "n": n, "age_sec": age}
|
|
|
|
def _job_counts(sid: str) -> Dict[str, int]:
|
|
outc: Dict[str, int] = {
|
|
"ls_N": 0, "ls_R": 0, "ls_O": 0,
|
|
"kw_I": 0, "kw_D": 0, "total": 0,
|
|
}
|
|
try:
|
|
rows = db.conn.execute(
|
|
"""
|
|
SELECT broker, job_flag, COUNT(*) AS n
|
|
FROM condition_job_events
|
|
WHERE strategy_id=%s AND event_time >= %s AND event_time <= %s
|
|
GROUP BY broker, job_flag
|
|
""",
|
|
(sid, start_s, end_s),
|
|
).fetchall()
|
|
except Exception:
|
|
return outc
|
|
for r in rows or []:
|
|
br = str(r["broker"] if hasattr(r, "keys") else r[0] or "").lower()
|
|
jf = str(r["job_flag"] if hasattr(r, "keys") else r[1] or "").upper()
|
|
n = int(r["n"] if hasattr(r, "keys") else r[2] or 0)
|
|
outc["total"] += n
|
|
key = None
|
|
if br == "ls" and jf in ("N", "R", "O"):
|
|
key = f"ls_{jf}"
|
|
elif br == "kiwoom" and jf in ("I", "D"):
|
|
key = f"kw_{jf}"
|
|
if key:
|
|
outc[key] = n
|
|
return outc
|
|
|
|
timeline: List[Dict[str, Any]] = []
|
|
try:
|
|
trows = db.conn.execute(
|
|
"""
|
|
SELECT LEFT(event_time, 16) AS bucket,
|
|
job_flag, COUNT(*) AS n
|
|
FROM condition_job_events
|
|
WHERE event_time >= %s AND event_time <= %s
|
|
GROUP BY LEFT(event_time, 16), job_flag
|
|
ORDER BY bucket
|
|
""",
|
|
(start_s, end_s),
|
|
).fetchall()
|
|
agg: Dict[str, Dict[str, int]] = {}
|
|
for r in trows or []:
|
|
raw = str(r["bucket"] if hasattr(r, "keys") else r[0] or "")
|
|
# 'YYYY-MM-DD HH:MM' → 10분 버킷
|
|
try:
|
|
base = raw[:14] # 'YYYY-MM-DD HH:'
|
|
minute = int(raw[14:16])
|
|
bkey = f"{base}{(minute // 10) * 10:02d}"
|
|
except Exception:
|
|
bkey = raw
|
|
jf = str(r["job_flag"] if hasattr(r, "keys") else r[1] or "").upper()
|
|
n = int(r["n"] if hasattr(r, "keys") else r[2] or 0)
|
|
slot = agg.setdefault(bkey, {"enter": 0, "exit": 0, "reenter": 0})
|
|
if jf in ("N", "I"):
|
|
slot["enter"] += n
|
|
elif jf == "R":
|
|
slot["reenter"] += n
|
|
elif jf in ("O", "D"):
|
|
slot["exit"] += n
|
|
timeline = [{"t": k, **v} for k, v in sorted(agg.items())]
|
|
except Exception as e:
|
|
logger.debug("parity timeline: %s", e)
|
|
|
|
strategies: List[Dict[str, Any]] = []
|
|
for sid in sids:
|
|
live_src = str(
|
|
get_env_from_db(f"{sid}_UNIVERSE_SOURCE", "") or ""
|
|
).strip().lower() or "—"
|
|
ls = _latest_hist("ls_candidates_history", sid)
|
|
kw = _latest_hist("target_candidates_history", sid)
|
|
jobs = _job_counts(sid)
|
|
status = "ok"
|
|
warn = ""
|
|
if live_src in ("ls_condition", "ls", "ls_afr"):
|
|
if (ls.get("n") or 0) <= 0:
|
|
status = "danger"
|
|
warn = "실매=LS 인데 당일 LS history 최신 스냅 0/없음 → 유니버스0 위험"
|
|
if (kw.get("n") or 0) > 0:
|
|
warn += f" (키움 history는 {kw['n']}종 있음 — t1859공백/AFR미적재 의심)"
|
|
elif ls.get("age_sec") is not None and ls["age_sec"] > 600:
|
|
status = "warn"
|
|
warn = f"LS history {ls['age_sec']}초 전 스냅 — 장중 공백 가능"
|
|
elif live_src in ("kiwoom_condition", "kiwoom"):
|
|
if (kw.get("n") or 0) <= 0:
|
|
status = "danger"
|
|
warn = "실매=키움 인데 당일 키움 history 없음"
|
|
elif kw.get("age_sec") is not None and kw["age_sec"] > 600:
|
|
status = "warn"
|
|
warn = f"키움 history {kw['age_sec']}초 전"
|
|
|
|
strategies.append({
|
|
"strategy_id": sid,
|
|
"live_source": live_src,
|
|
"ls_hist_n": ls.get("n") or 0,
|
|
"ls_hist_at": ls.get("event_time"),
|
|
"ls_age_sec": ls.get("age_sec"),
|
|
"kw_hist_n": kw.get("n") or 0,
|
|
"kw_hist_at": kw.get("event_time"),
|
|
"kw_age_sec": kw.get("age_sec"),
|
|
"jobs": jobs,
|
|
"status": status,
|
|
"warn": warn,
|
|
})
|
|
|
|
try:
|
|
grace_sec = int(str(get_env_from_db("CONDITION_EXIT_GRACE_SEC", "0") or "0"))
|
|
except (TypeError, ValueError):
|
|
grace_sec = 0
|
|
|
|
return {
|
|
"day": day_s,
|
|
"grace_sec": grace_sec,
|
|
"strategies": strategies,
|
|
"timeline": timeline,
|
|
"notes": (
|
|
"history=봇이 저장한 effective 스냅. 실매 후보는 키움 RAM. "
|
|
"실매 소스가 LS인데 LS history 0 이면 universe_zero 알림이 정상"
|
|
"(키움 이중이력과 무관)."
|
|
),
|
|
}
|